diff --git a/tests/kicad/utils/occ-service.ts b/tests/kicad/utils/occ-service.ts index bdc027a..46bcbb7 100644 --- a/tests/kicad/utils/occ-service.ts +++ b/tests/kicad/utils/occ-service.ts @@ -45,58 +45,164 @@ export async function installOccServiceStub(page: Page): Promise { (window as any).__occExports = []; - let workerP: Promise | null = null; - const pending = new Map void>(); - let nextId = 1; + interface WorkerSlot { + generation: number; + worker?: Worker; + failed: boolean; + ready: Promise; + /** The exact lifecycle transition used by Worker.onmessageerror. */ + failDecode: () => void; + } - const ensureWorker = (): Promise => { - if (!workerP) { - workerP = (async () => { + interface PendingRequest { + generation: number; + resolve: (res: any) => void; + } + + let nextGeneration = 1; + let workerSlot: WorkerSlot | null = null; + const pending = new Map(); + let nextId = 1; + let maxPending = 0; + let requestsStarted = 0; + let requestsPosted = 0; + const workerGenerationsStarted: number[] = []; + let armedFault: { + count: number; + kind: 'terminate' | 'messageerror'; + report: string; + } | null = null; + const retiredGenerations: number[] = []; + + const pendingInGeneration = (generation: number): number => { + let count = 0; + for (const request of pending.values()) { + if (request.generation === generation) count++; + } + return count; + }; + + const failPending = (generation: number, report: string): void => { + for (const [id, request] of pending) { + if (request.generation !== generation) continue; + pending.delete(id); + request.resolve({ ok: false, report }); + } + }; + + const retireWorker = (slot: WorkerSlot, report: string): void => { + if (slot.failed) return; + slot.failed = true; + retiredGenerations.push(slot.generation); + failPending(slot.generation, report); + if (workerSlot === slot) workerSlot = null; + try { + slot.worker?.terminate(); + } catch { + /* already gone */ + } + }; + + const maybeTriggerArmedFault = (slot: WorkerSlot): void => { + if (!armedFault || slot.failed) return; + if (pendingInGeneration(slot.generation) < armedFault.count) return; + const { kind, report } = armedFault; + armedFault = null; + console.log(`[TEST-OCC] faulting generation ${slot.generation} (${kind}): ${report}`); + if (kind === 'messageerror' && slot.worker) { + // Synthetic dispatch on Worker is engine-dependent. Invoke + // the exact transition installed as the real event handler. + slot.failDecode(); + } else { + // Worker.terminate() intentionally emits no error event, so + // the hook supplies the fatal lifecycle transition explicitly. + retireWorker(slot, report); + } + }; + + const ensureWorker = (): Promise => { + if (!workerSlot) { + const slot = { + generation: nextGeneration++, + failed: false, + } as WorkerSlot; + workerGenerationsStarted.push(slot.generation); + + slot.ready = (async () => { const glue = new URL('occ_service.js', window.location.href).href; console.log(`[TEST-OCC] booting occ_service from ${glue}`); const worker = new Worker(URL.createObjectURL(new Blob( [`self.OCC_GLUE_URL = ${JSON.stringify(glue)};\n`, workerSrc], { type: 'text/javascript' }))); + slot.worker = worker; + let rejectBoot: ((reason?: unknown) => void) | undefined; + let removeBootListener: (() => void) | undefined; + + worker.onerror = (e) => { + const report = `occ_service crashed: ${e.message || 'worker error'}`; + console.error(`[TEST-OCC] ${report}; resetting service`); + retireWorker(slot, report); + removeBootListener?.(); + removeBootListener = undefined; + const reject = rejectBoot; + rejectBoot = undefined; + reject?.(new Error(report)); + }; + slot.failDecode = () => { + const report = 'occ_service transport failed: message decode failed'; + console.error(`[TEST-OCC] ${report}; resetting service`); + retireWorker(slot, report); + removeBootListener?.(); + removeBootListener = undefined; + const reject = rejectBoot; + rejectBoot = undefined; + reject?.(new Error(report)); + }; + worker.onmessageerror = slot.failDecode; worker.onmessage = (e) => { + if (slot.failed || workerSlot !== slot) return; const { id, res } = e.data ?? {}; if (typeof id !== 'number') return; - const resolve = pending.get(id); - if (resolve) { pending.delete(id); resolve(res); } + const request = pending.get(id); + if (request?.generation === slot.generation) { + pending.delete(id); + request.resolve(res); + } }; - // Legible boot: the old handshake could never reject on a - // worker DEATH (importScripts throw, pthread spawn wedge, - // OOM-kill) — the promise just hung until the spec's 180s - // timeout with zero evidence. Surface worker errors and - // bound the boot. await new Promise((resolve, reject) => { - const fail = (msg: string) => { - clearTimeout(timer); - reject(new Error(msg)); - }; - const timer = setTimeout( - () => fail('[TEST-OCC] occ_service boot timed out after 60s ' - + '(no ready/bootError from the worker)'), 60000); + rejectBoot = reject; const onFirst = (e: MessageEvent) => { if (e.data?.ready) { - worker.removeEventListener('message', onFirst); - clearTimeout(timer); + removeBootListener?.(); + removeBootListener = undefined; + rejectBoot = undefined; resolve(); } else if (e.data?.bootError) { - fail(`[TEST-OCC] occ_service bootError: ${e.data.bootError}`); + removeBootListener?.(); + removeBootListener = undefined; + rejectBoot = undefined; + const report = `occ_service boot failed: ${String(e.data.bootError)}`; + retireWorker(slot, report); + reject(new Error(report)); } }; worker.addEventListener('message', onFirst); - worker.addEventListener('error', (e: any) => fail( - `[TEST-OCC] occ_service worker error: ${e?.message ?? e} ` - + `(${e?.filename ?? '?'}:${e?.lineno ?? '?'})`)); - worker.addEventListener('messageerror', () => fail( - '[TEST-OCC] occ_service worker messageerror (structured clone failed)')); + removeBootListener = () => worker.removeEventListener('message', onFirst); }); + if (slot.failed || workerSlot !== slot) + throw new Error('occ_service worker retired during boot'); console.log('[TEST-OCC] occ_service ready'); - return worker; - })().catch((e) => { workerP = null; throw e; }); + return slot; + })().catch((e) => { + // A late rejection from a retired generation cannot clear + // the replacement slot created by a new request. + retireWorker(slot, `occ_service unavailable: ${String(e)}`); + throw e; + }); + + workerSlot = slot; } - return workerP; + return workerSlot.ready; }; // Mirror of the app's collectBoardModelFiles, against the page's @@ -127,21 +233,36 @@ export async function installOccServiceStub(page: Page): Promise { }; const request = async (req: any) => { + // Count provider entry before model collection or worker boot. This + // distinguishes "the wx button reached OCC" from a worker that was + // already active for some earlier request. + requestsStarted++; if (req.kind === 'export') req.models = await collectModels(new TextDecoder().decode(req.board)); - let worker: Worker; + let slot: WorkerSlot; try { - worker = await ensureWorker(); + slot = await ensureWorker(); } catch (e) { return { ok: false, report: `occ_service unavailable: ${e}` }; } + const worker = slot.worker; + if (!worker || slot.failed || workerSlot !== slot) + return { ok: false, report: 'occ_service worker is unavailable' }; const id = nextId++; const transfer = req.kind === 'export' ? [req.board.buffer, ...(req.models ?? []).map((m: any) => m.bytes.buffer)] : [req.bytes.buffer]; const res: any = await new Promise((resolve) => { - pending.set(id, resolve); - worker.postMessage({ id, req }, transfer); + pending.set(id, { generation: slot.generation, resolve }); + maxPending = Math.max(maxPending, pendingInGeneration(slot.generation)); + try { + worker.postMessage({ id, req }, transfer); + requestsPosted++; + maybeTriggerArmedFault(slot); + } catch (error) { + pending.delete(id); + resolve({ ok: false, report: `occ_service request failed: ${String(error)}` }); + } }); if (req.kind === 'export') { if (res.ok && res.bytes?.length) { @@ -169,6 +290,39 @@ export async function installOccServiceStub(page: Page): Promise { return res; }; + (globalThis as any).__occServiceTestHooks = { + /** Arm one deterministic host-side fault after N real posts. */ + terminateWhenPendingAtLeast(count: number, report = 'occ_service test fault') { + if (!Number.isSafeInteger(count) || count < 1) + throw new Error('pending threshold must be a positive safe integer'); + armedFault = { count, kind: 'terminate', report }; + if (workerSlot) maybeTriggerArmedFault(workerSlot); + }, + /** Arm the real Worker's production-parity messageerror handler. */ + messageErrorWhenPendingAtLeast(count: number) { + if (!Number.isSafeInteger(count) || count < 1) + throw new Error('pending threshold must be a positive safe integer'); + armedFault = { + count, + kind: 'messageerror', + report: 'occ_service transport failed: message decode failed', + }; + if (workerSlot) maybeTriggerArmedFault(workerSlot); + }, + snapshot() { + return { + activeGeneration: workerSlot?.generation ?? null, + pending: pending.size, + maxPending, + requestsStarted, + requestsPosted, + workerGenerationsStarted: [...workerGenerationsStarted], + retiredGenerations: [...retiredGenerations], + armed: armedFault !== null, + }; + }, + }; + (globalThis as any).occService = { request }; }, OCC_WORKER_SRC); } diff --git a/web/standalone/src/wasm/libs/models-bridge.test.ts b/web/standalone/src/wasm/libs/models-bridge.test.ts index 685e3de..e4f90f5 100644 --- a/web/standalone/src/wasm/libs/models-bridge.test.ts +++ b/web/standalone/src/wasm/libs/models-bridge.test.ts @@ -1,4 +1,4 @@ -import { describe, expect, it } from "vitest"; +import { describe, expect, it, vi } from "vitest"; import { collectBoardModelFiles, ensureModelInMemfs, @@ -140,6 +140,64 @@ describe("collectBoardModelFiles", () => { installFakes(() => true); expect(await collectBoardModelFiles("(kicad_pcb (version 1))")).toEqual([]); }); + + it("makes an in-flight source result inert and starts no later ref after abort", async () => { + // repro for E-4: the prefetch abort must stop selection immediately and + // may not retain a body that resolves after the abort. + (globalThis as unknown as { window: unknown }).window ??= globalThis; + let resolveFirst!: (body: Uint8Array | null) => void; + const source: Model3dSource = { + getModelBody: vi.fn( + () => new Promise((resolve) => { + resolveFirst = resolve; + }), + ), + hasModel: async () => true, + }; + installModel3dHandler(source, () => {}); + const controller = new AbortController(); + const retired = new Error("exact OCC prefetch retired"); + const collection = collectBoardModelFiles( + '(model "AbortA.3dshapes/A.step")\n' + + '(model "AbortB.3dshapes/B.step")', + 1, + controller.signal, + ); + + await vi.waitFor(() => expect(source.getModelBody).toHaveBeenCalledTimes(1)); + controller.abort(retired); + resolveFirst(new Uint8Array([1, 2, 3])); + + await expect(collection).rejects.toBe(retired); + expect(source.getModelBody).toHaveBeenCalledTimes(1); + }); + + it("never touches the editor MEMFS (pure source/IDB/network path)", async () => { + // repro for E-4: routing bodies through the editor heap added a stale + // native-completion tail after a prefetch timeout — the collect path must + // stay off FS entirely (the OCC worker stages bodies in its own MEMFS). + const fs = { + mkdirTree: vi.fn(), + writeFile: vi.fn(), + analyzePath: vi.fn(() => ({ exists: false })), + readFile: vi.fn(), + }; + (globalThis as unknown as { window: unknown }).window ??= globalThis; + (globalThis as unknown as { FS: unknown }).FS = fs; + const source: Model3dSource = { + getModelBody: async (ref) => new TextEncoder().encode(`body:${ref}`), + hasModel: async () => true, + }; + installModel3dHandler(source, () => {}); + const models = await collectBoardModelFiles( + '(model "PureLib.3dshapes/M1.step")', + ); + expect(models).toHaveLength(1); + expect(fs.mkdirTree).not.toHaveBeenCalled(); + expect(fs.writeFile).not.toHaveBeenCalled(); + expect(fs.readFile).not.toHaveBeenCalled(); + expect(fs.analyzePath).not.toHaveBeenCalled(); + }); }); describe("scanModelRefs", () => { diff --git a/web/standalone/src/wasm/libs/models-bridge.ts b/web/standalone/src/wasm/libs/models-bridge.ts index 742bfdd..16585e3 100644 --- a/web/standalone/src/wasm/libs/models-bridge.ts +++ b/web/standalone/src/wasm/libs/models-bridge.ts @@ -195,21 +195,29 @@ export interface BoardModelFile { } /** - * Prefetch + read back every lib model a board references, for shipping with - * an occ_service export request — the worker is its own wasm module with its - * own MEMFS, so the editor-side files are invisible there. Reuses - * ensureModelInMemfs (IDB/R2-cached, coalesced, wrl→step format fallback); - * the returned paths carry the staged file's REAL extension, deduplicated - * (two refs can materialize to the same substituted body). Best-effort: a - * ref the source can't serve is skipped (the exporter reports it missing). - * Returns [] when 3D model delivery is not configured. + * Fetch every lib model a board references for an occ_service export. This is + * deliberately a pure source/IDB/network path (E-4): the OCC worker has a + * different MEMFS, so materializing and reading the bytes through the editor's + * native heap only adds a stale native-completion tail after an + * export-prefetch timeout. The returned paths carry the fetched file's real + * fallback extension and are deduplicated. Best-effort: missing models are + * skipped (the exporter reports them missing). An abort stops selection + * immediately and makes every already-started source result inert before it + * is retained. Returns [] when 3D model delivery is not configured. */ export async function collectBoardModelFiles( boardText: string, concurrency = 6, + signal?: AbortSignal, ): Promise { - const fs = toolFS(); - if (!installedSource || !fs) return []; + if (!installedSource) return []; + const source = installedSource; + const isCurrent = () => source === installedSource; + const throwIfAborted = (): void => { + if (!signal?.aborted) return; + throw signal.reason ?? new DOMException("Model collection aborted", "AbortError"); + }; + throwIfAborted(); const refs = scanModelRefs(boardText); if (!refs.length) return []; @@ -217,25 +225,39 @@ export async function collectBoardModelFiles( const seen = new Set(); let idx = 0; const worker = async (): Promise => { - while (idx < refs.length) { + while (true) { + throwIfAborted(); + if (!isCurrent() || idx >= refs.length) return; const ref = refs[idx++]!; - try { - const abs = await ensureModelInMemfs(ref); - if (!abs || seen.has(abs)) continue; - seen.add(abs); - // FS.readFile copies out of the wasm heap — the buffer is safely - // transferable to the worker. - const bytes = fs.readFile(abs) as Uint8Array; - out.push({ path: abs.slice(MODELS_3D_ROOT.length + 1), bytes }); - } catch { - // best-effort: a missing body surfaces as the exporter's own - // "Could not add 3D model" report warning, never a failed export + for (const candidate of refCandidates(ref)) { + let body: Uint8Array | null = null; + try { + body = await source.getModelBody(candidate); + } catch { + // Best-effort: try the next format. The exporter reports a miss if + // no candidate exists. + } + throwIfAborted(); + if (!isCurrent()) return; + if (!body) continue; + + if (!seen.has(candidate)) { + seen.add(candidate); + // The OCC worker receives this buffer as a transferable. Keep its + // ownership independent from any source/cache view. + out.push({ path: candidate, bytes: new Uint8Array(body) }); + } + break; } } }; await Promise.all( - Array.from({ length: Math.min(concurrency, refs.length) }, () => worker()), + Array.from( + { length: Math.min(Math.max(1, Math.trunc(concurrency)), refs.length) }, + () => worker(), + ), ); + throwIfAborted(); installedLog(`[3d] export prefetch: ${out.length}/${refs.length} board model(s)`); return out; } diff --git a/web/standalone/src/wasm/occ-service.test.ts b/web/standalone/src/wasm/occ-service.test.ts new file mode 100644 index 0000000..19a74d7 --- /dev/null +++ b/web/standalone/src/wasm/occ-service.test.ts @@ -0,0 +1,413 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +vi.mock("@/lib/download", () => ({ downloadBytes: vi.fn() })); +vi.mock("./libs/models-bridge", () => ({ + collectBoardModelFiles: vi.fn(async () => []), +})); +vi.mock("./occ-worker.js?raw", () => ({ default: "// fake OCC worker" })); +vi.mock("./wasm-assets", () => ({ + resolveWasmBase: vi.fn(async () => "/wasm"), +})); + +import { installOccService, type OccResponse } from "./occ-service"; +import { collectBoardModelFiles } from "./libs/models-bridge"; +import { resolveWasmBase } from "./wasm-assets"; + +const mockedCollectBoardModelFiles = vi.mocked(collectBoardModelFiles); +const mockedResolveWasmBase = vi.mocked(resolveWasmBase); + +const TEST_MODEL_PREFETCH_TIMEOUT_MS = 500; +const TEST_BOOT_TIMEOUT_MS = 1_000; +const TEST_RESPONSE_TIMEOUT_MS = 5_000; + +type MessageListener = (event: MessageEvent) => void; + +class FakeWorker { + static instances: FakeWorker[] = []; + + onmessage: MessageListener | null = null; + onerror: ((event: ErrorEvent) => void) | null = null; + onmessageerror: ((event: MessageEvent) => void) | null = null; + readonly postMessage = vi.fn(); + readonly terminate = vi.fn(); + private readonly messageListeners = new Set(); + + constructor() { + FakeWorker.instances.push(this); + } + + addEventListener(type: string, listener: EventListenerOrEventListenerObject): void { + if (type === "message") this.messageListeners.add(listener as MessageListener); + } + + removeEventListener(type: string, listener: EventListenerOrEventListenerObject): void { + if (type === "message") this.messageListeners.delete(listener as MessageListener); + } + + emitMessage(data: unknown): void { + const event = { data } as MessageEvent; + this.onmessage?.(event); + for (const listener of [...this.messageListeners]) listener(event); + } + + emitError(message: string): void { + this.onerror?.({ message } as ErrorEvent); + } + + emitMessageError(): void { + this.onmessageerror?.({} as MessageEvent); + } +} + +const loadRequest = () => ({ + kind: "loadModel" as const, + bytes: new Uint8Array([1, 2, 3]), + ext: "step", +}); + +const exportRequest = () => ({ + kind: "export" as const, + board: new TextEncoder().encode("(kicad_pcb)"), + jobJson: '{"type":"step"}', + fileName: "board.step", +}); + +const service = () => { + const installed = globalThis.occService; + if (!installed) throw new Error("OCC service was not installed"); + return installed; +}; + +async function waitForWorker(index: number): Promise { + await vi.waitFor(() => expect(FakeWorker.instances.length).toBeGreaterThan(index)); + return FakeWorker.instances[index]!; +} + +async function readyRequest(workerIndex: number): Promise<{ + worker: FakeWorker; + request: Promise; + id: number; +}> { + const request = service().request(loadRequest()); + const worker = await waitForWorker(workerIndex); + worker.emitMessage({ ready: true }); + await vi.waitFor(() => expect(worker.postMessage).toHaveBeenCalledTimes(1)); + const [{ id }] = worker.postMessage.mock.calls[0] as [{ id: number }]; + return { worker, request, id }; +} + +describe("OCC service worker lifetime", () => { + beforeEach(() => { + FakeWorker.instances = []; + mockedCollectBoardModelFiles.mockReset(); + mockedCollectBoardModelFiles.mockResolvedValue([]); + mockedResolveWasmBase.mockReset(); + mockedResolveWasmBase.mockResolvedValue("/wasm"); + vi.stubGlobal("window", { location: { href: "https://pcbjam.test/editor" } }); + vi.stubGlobal("Worker", FakeWorker); + let nextBlob = 1; + vi.spyOn(URL, "createObjectURL").mockImplementation( + () => `blob:occ-test-${nextBlob++}`, + ); + vi.spyOn(URL, "revokeObjectURL").mockImplementation(() => undefined); + delete globalThis.occService; + installOccService(vi.fn(), { + modelPrefetchTimeoutMs: TEST_MODEL_PREFETCH_TIMEOUT_MS, + bootTimeoutMs: TEST_BOOT_TIMEOUT_MS, + responseTimeoutMs: TEST_RESPONSE_TIMEOUT_MS, + }); + }); + + afterEach(() => { + delete globalThis.occService; + vi.useRealTimers(); + vi.unstubAllGlobals(); + vi.restoreAllMocks(); + }); + + it("bounds a never-ready generation, retires it, and boots a fresh worker", async () => { + vi.useFakeTimers(); + + const timedOut = service().request(loadRequest()); + await vi.advanceTimersByTimeAsync(0); + const deadWorker = FakeWorker.instances[0]!; + expect(deadWorker).toBeDefined(); + + await vi.advanceTimersByTimeAsync(TEST_BOOT_TIMEOUT_MS); + await expect(timedOut).resolves.toMatchObject({ + ok: false, + report: expect.stringContaining( + `occ_service boot timed out after ${TEST_BOOT_TIMEOUT_MS} ms`, + ), + }); + expect(deadWorker.terminate).toHaveBeenCalledTimes(1); + expect(URL.revokeObjectURL).toHaveBeenCalledWith("blob:occ-test-1"); + expect(vi.getTimerCount()).toBe(0); + + deadWorker.emitMessage({ ready: true }); + + const recovered = service().request(loadRequest()); + await vi.advanceTimersByTimeAsync(0); + const freshWorker = FakeWorker.instances[1]!; + expect(freshWorker).toBeDefined(); + freshWorker.emitMessage({ ready: true }); + await vi.advanceTimersByTimeAsync(0); + const [{ id }] = freshWorker.postMessage.mock.calls[0] as [{ id: number }]; + freshWorker.emitMessage({ id, res: { ok: true, report: "fresh" } }); + + await expect(recovered).resolves.toEqual({ ok: true, report: "fresh" }); + expect(freshWorker.terminate).not.toHaveBeenCalled(); + expect(vi.getTimerCount()).toBe(0); + }); + + it("bounds delivery resolution and never creates a late retired worker", async () => { + vi.useFakeTimers(); + let releaseBase!: (base: string) => void; + mockedResolveWasmBase.mockImplementationOnce( + () => new Promise((resolve) => { + releaseBase = resolve; + }), + ); + + const timedOut = service().request(loadRequest()); + await vi.advanceTimersByTimeAsync(0); + expect(FakeWorker.instances).toHaveLength(0); + + await vi.advanceTimersByTimeAsync(TEST_BOOT_TIMEOUT_MS); + await expect(timedOut).resolves.toMatchObject({ + ok: false, + report: expect.stringContaining( + `occ_service boot timed out after ${TEST_BOOT_TIMEOUT_MS} ms`, + ), + }); + expect(FakeWorker.instances).toHaveLength(0); + expect(vi.getTimerCount()).toBe(0); + + // Delivery may still finish because the underlying fetch is not abortable + // here. Its retired generation must not create a Worker or replace the + // fresh slot which the next exact request owns. + releaseBase("/stale-wasm"); + await vi.advanceTimersByTimeAsync(0); + expect(FakeWorker.instances).toHaveLength(0); + + const recovered = service().request(loadRequest()); + await vi.advanceTimersByTimeAsync(0); + const freshWorker = FakeWorker.instances[0]!; + expect(freshWorker).toBeDefined(); + freshWorker.emitMessage({ ready: true }); + await vi.advanceTimersByTimeAsync(0); + const [{ id }] = freshWorker.postMessage.mock.calls[0] as [{ id: number }]; + freshWorker.emitMessage({ id, res: { ok: true, report: "fresh" } }); + + await expect(recovered).resolves.toEqual({ ok: true, report: "fresh" }); + expect(vi.getTimerCount()).toBe(0); + }); + + it("exports without a hung model prefetch and ignores its late result", async () => { + vi.useFakeTimers(); + let releaseModels!: ( + models: Array<{ path: string; bytes: Uint8Array }>, + ) => void; + let prefetchSignal: AbortSignal | undefined; + mockedCollectBoardModelFiles.mockImplementationOnce( + (_board, _concurrency, signal) => { + prefetchSignal = signal; + return new Promise((resolve) => { + releaseModels = resolve; + }); + }, + ); + const input = exportRequest(); + const request = service().request(input); + + // Request fields are captured before optional asynchronous preparation. + input.fileName = "mutated-after-dispatch.step"; + await vi.advanceTimersByTimeAsync(0); + expect(FakeWorker.instances).toHaveLength(0); + + await vi.advanceTimersByTimeAsync(TEST_MODEL_PREFETCH_TIMEOUT_MS); + expect(prefetchSignal?.aborted).toBe(true); + const worker = FakeWorker.instances[0]!; + expect(worker).toBeDefined(); + worker.emitMessage({ ready: true }); + await vi.advanceTimersByTimeAsync(0); + expect(worker.postMessage).toHaveBeenCalledTimes(1); + const [{ id, req: dispatched }] = worker.postMessage.mock.calls[0] as [ + { + id: number; + req: { + kind: "export"; + fileName: string; + models: Array<{ path: string; bytes: Uint8Array }>; + }; + }, + ]; + expect(dispatched).not.toBe(input); + expect(dispatched.fileName).toBe("board.step"); + expect(dispatched.models).toEqual([]); + expect(input).not.toHaveProperty("models"); + + releaseModels([ + { path: "Late.3dshapes/model.step", bytes: new Uint8Array([9]) }, + ]); + await vi.advanceTimersByTimeAsync(0); + expect(dispatched.models).toEqual([]); + expect(worker.postMessage).toHaveBeenCalledTimes(1); + + worker.emitMessage({ id, res: { ok: true, report: "exported" } }); + await expect(request).resolves.toEqual({ + ok: true, + report: "exported", + fileName: undefined, + }); + expect(vi.getTimerCount()).toBe(0); + }); + + it("retires a ready-but-silent generation and settles all concurrent ids", async () => { + vi.useFakeTimers(); + + const first = service().request(loadRequest()); + await vi.advanceTimersByTimeAsync(0); + const deadWorker = FakeWorker.instances[0]!; + deadWorker.emitMessage({ ready: true }); + await vi.advanceTimersByTimeAsync(0); + + const second = service().request(loadRequest()); + await vi.advanceTimersByTimeAsync(0); + expect(deadWorker.postMessage).toHaveBeenCalledTimes(2); + const [{ id: firstId }] = deadWorker.postMessage.mock.calls[0] as [ + { id: number }, + ]; + const [{ id: secondId }] = deadWorker.postMessage.mock.calls[1] as [ + { id: number }, + ]; + expect(firstId).not.toBe(secondId); + + await vi.advanceTimersByTimeAsync(TEST_RESPONSE_TIMEOUT_MS); + const timeout = + `occ_service response timed out after ${TEST_RESPONSE_TIMEOUT_MS} ms`; + await expect(first).resolves.toEqual({ ok: false, report: timeout }); + await expect(second).resolves.toEqual({ ok: false, report: timeout }); + expect(deadWorker.terminate).toHaveBeenCalledTimes(1); + expect(URL.revokeObjectURL).toHaveBeenCalledWith("blob:occ-test-1"); + expect(vi.getTimerCount()).toBe(0); + + const recovered = service().request(loadRequest()); + await vi.advanceTimersByTimeAsync(0); + const freshWorker = FakeWorker.instances[1]!; + freshWorker.emitMessage({ ready: true }); + await vi.advanceTimersByTimeAsync(0); + const [{ id: freshId }] = freshWorker.postMessage.mock.calls[0] as [ + { id: number }, + ]; + let recoveredSettled = false; + const observed = recovered.then((response) => { + recoveredSettled = true; + return response; + }); + + deadWorker.emitMessage({ + id: freshId, + res: { ok: true, report: "stale" }, + }); + await Promise.resolve(); + expect(recoveredSettled).toBe(false); + expect(freshWorker.terminate).not.toHaveBeenCalled(); + + freshWorker.emitMessage({ + id: freshId, + res: { ok: true, report: "fresh" }, + }); + await expect(observed).resolves.toEqual({ ok: true, report: "fresh" }); + expect(vi.getTimerCount()).toBe(0); + }); + + it("settles every pending request on a runtime crash and restarts", async () => { + const first = await readyRequest(0); + const alsoPending = service().request(loadRequest()); + await vi.waitFor(() => expect(first.worker.postMessage).toHaveBeenCalledTimes(2)); + + first.worker.emitError("wasm trap"); + + await expect(first.request).resolves.toEqual({ + ok: false, + report: "occ_service crashed: wasm trap", + }); + await expect(alsoPending).resolves.toEqual({ + ok: false, + report: "occ_service crashed: wasm trap", + }); + expect(first.worker.terminate).toHaveBeenCalledTimes(1); + + const second = await readyRequest(1); + expect(second.worker).not.toBe(first.worker); + + let secondSettled = false; + const observedSecond = second.request.then((response) => { + secondSettled = true; + return response; + }); + + // A callback retained by the retired generation cannot consume the current + // generation's request, even if it carries that request's numeric id. + first.worker.emitMessage({ + id: second.id, + res: { ok: true, report: "stale result" }, + }); + await Promise.resolve(); + expect(secondSettled).toBe(false); + + second.worker.emitMessage({ + id: second.id, + res: { ok: true, report: "fresh result" }, + }); + await expect(observedSecond).resolves.toEqual({ + ok: true, + report: "fresh result", + }); + }); + + it("turns a boot error into a response and keeps the next boot retryable", async () => { + const failedRequest = service().request(loadRequest()); + const failedWorker = await waitForWorker(0); + + failedWorker.emitMessage({ bootError: "initialization failed" }); + + await expect(failedRequest).resolves.toMatchObject({ + ok: false, + report: expect.stringContaining("occ_service boot failed: initialization failed"), + }); + expect(failedWorker.terminate).toHaveBeenCalledTimes(1); + + const retry = await readyRequest(1); + retry.worker.emitMessage({ + id: retry.id, + res: { ok: true, report: "recovered" }, + }); + await expect(retry.request).resolves.toEqual({ ok: true, report: "recovered" }); + }); + + it("fails every request in a decode-faulted generation, then recovers", async () => { + const failed = await readyRequest(0); + const alsoPending = service().request(loadRequest()); + await vi.waitFor(() => expect(failed.worker.postMessage).toHaveBeenCalledTimes(2)); + failed.worker.emitMessageError(); + + await expect(failed.request).resolves.toEqual({ + ok: false, + report: "occ_service transport failed: message decode failed", + }); + await expect(alsoPending).resolves.toEqual({ + ok: false, + report: "occ_service transport failed: message decode failed", + }); + expect(failed.worker.terminate).toHaveBeenCalledTimes(1); + + const retry = await readyRequest(1); + retry.worker.emitMessage({ + id: retry.id, + res: { ok: true, report: "decoded" }, + }); + await expect(retry.request).resolves.toEqual({ ok: true, report: "decoded" }); + }); +}); diff --git a/web/standalone/src/wasm/occ-service.ts b/web/standalone/src/wasm/occ-service.ts index 319e3cc..2436496 100644 --- a/web/standalone/src/wasm/occ-service.ts +++ b/web/standalone/src/wasm/occ-service.ts @@ -67,101 +67,314 @@ export function occWorkerBlobParts(glueHref: string): string[] { ]; } -export function installOccService(log: (msg: string) => void): void { +export interface OccServiceWatchdogs { + /** Maximum time to wait for optional board-model prefetch before exporting without it. */ + modelPrefetchTimeoutMs?: number; + /** Maximum time from the first request until a new generation announces `ready`. */ + bootTimeoutMs?: number; + /** Maximum time for any one request in a ready generation to answer. */ + responseTimeoutMs?: number; +} + +// These are last-resort failure bounds, not normal scheduling deadlines. +// OCC startup, model parsing, and board export can all be expensive on slow +// devices, so production defaults deliberately leave a large margin. +export const OCC_BOOT_TIMEOUT_MS = 2 * 60_000; +export const OCC_RESPONSE_TIMEOUT_MS = 30 * 60_000; +export const OCC_MODEL_PREFETCH_TIMEOUT_MS = 30_000; + +export function installOccService( + log: (msg: string) => void, + watchdogs: OccServiceWatchdogs = {}, +): void { if (globalThis.occService) return; - let nextId = 1; - const pending = new Map void>(); - let workerP: Promise | null = null; + const modelPrefetchTimeoutMs = + watchdogs.modelPrefetchTimeoutMs ?? OCC_MODEL_PREFETCH_TIMEOUT_MS; + const bootTimeoutMs = watchdogs.bootTimeoutMs ?? OCC_BOOT_TIMEOUT_MS; + const responseTimeoutMs = + watchdogs.responseTimeoutMs ?? OCC_RESPONSE_TIMEOUT_MS; - const ensureWorker = (): Promise => { - if (!workerP) { - workerP = (async () => { + interface WorkerSlot { + generation: number; + worker?: Worker; + workerUrl?: string; + failed: boolean; + ready: Promise; + bootTimer?: ReturnType; + rejectBoot?: (reason?: unknown) => void; + removeBootListener?: () => void; + } + + interface PendingRequest { + generation: number; + resolve: (res: OccResponse) => void; + timer: ReturnType; + } + + let nextId = 1; + let nextGeneration = 1; + const pending = new Map(); + let workerSlot: WorkerSlot | null = null; + + const failPending = (generation: number, report: string): void => { + for (const [id, request] of pending) { + if (request.generation !== generation) continue; + pending.delete(id); + clearTimeout(request.timer); + request.resolve({ ok: false, report }); + } + }; + + const retireWorker = (slot: WorkerSlot, report: string): void => { + if (slot.failed) return; + slot.failed = true; + if (slot.bootTimer !== undefined) { + clearTimeout(slot.bootTimer); + slot.bootTimer = undefined; + } + slot.removeBootListener?.(); + slot.removeBootListener = undefined; + failPending(slot.generation, report); + if (workerSlot === slot) workerSlot = null; + try { + slot.worker?.terminate(); + } catch { + /* already gone */ + } + if (slot.workerUrl) { + try { + URL.revokeObjectURL(slot.workerUrl); + } catch { + /* URL cleanup must not prevent exact wait settlement */ + } + slot.workerUrl = undefined; + } + const reject = slot.rejectBoot; + slot.rejectBoot = undefined; + reject?.(new Error(report)); + }; + + const ensureWorker = (): Promise => { + if (!workerSlot) { + const slot = { + generation: nextGeneration++, + failed: false, + } as WorkerSlot; + // Publish the generation before its async boot reaches the first await. + // This also lets every continuation test exact slot ownership directly. + workerSlot = slot; + + // The editor is parked for this entire operation, including delivery + // discovery. Start the generation deadline before resolveWasmBase(): a + // hung manifest/CDN lookup must settle the exact wait just like a Worker + // which never announces ready. + const bootDeadline = new Promise((_resolve, reject) => { + slot.rejectBoot = reject; + slot.bootTimer = setTimeout(() => { + if (slot.failed || workerSlot !== slot) return; + const report = + `occ_service boot timed out after ${bootTimeoutMs} ms`; + log(`[occ] ${report} — resetting service`); + retireWorker(slot, report); + }, bootTimeoutMs); + }); + + const boot = (async () => { // occ_service is a Bundle (a published delivery artifact), not a Tool — // resolveWasmBase accepts either and looks the bundle up directly. const base = await resolveWasmBase("occ_service"); + if (slot.failed || workerSlot !== slot) { + throw new Error("occ_service worker retired during delivery resolution"); + } const glue = new URL(`${base}/occ_service.js`, window.location.href).href; log(`[occ] booting occ_service from ${base}`); - const worker = new Worker( - URL.createObjectURL( - new Blob(occWorkerBlobParts(glue), { type: "text/javascript" }), - ), + slot.workerUrl = URL.createObjectURL( + new Blob(occWorkerBlobParts(glue), { type: "text/javascript" }), ); + const worker = new Worker(slot.workerUrl); + slot.worker = worker; + + // A hard OCC/Wasm fault must complete every exact editor wait which + // depends on this worker. The next request gets a fresh generation; + // callbacks from this retired worker cannot resolve its requests. + worker.onerror = (e) => { + const report = `occ_service crashed: ${e.message || "worker error"}`; + log(`[occ] worker error: ${e.message || "worker error"} — resetting service`); + retireWorker(slot, report); + }; + worker.onmessageerror = () => { + const report = "occ_service transport failed: message decode failed"; + log("[occ] worker message decode failed — resetting service"); + retireWorker(slot, report); + }; worker.onmessage = (e) => { + if (slot.failed || workerSlot !== slot) return; const { id, res } = e.data ?? {}; if (typeof id !== "number") return; - const resolve = pending.get(id); - if (resolve) { + const request = pending.get(id); + if (request?.generation === slot.generation) { pending.delete(id); - resolve(res as OccResponse); + clearTimeout(request.timer); + request.resolve(res as OccResponse); } }; - await new Promise((resolve, reject) => { + await new Promise((resolve) => { const onFirst = (e: MessageEvent) => { if (e.data?.ready) { - worker.removeEventListener("message", onFirst); + if (slot.bootTimer !== undefined) { + clearTimeout(slot.bootTimer); + slot.bootTimer = undefined; + } + slot.removeBootListener?.(); + slot.removeBootListener = undefined; + slot.rejectBoot = undefined; resolve(); } else if (e.data?.bootError) { - reject(new Error(e.data.bootError)); + const report = `occ_service boot failed: ${String(e.data.bootError)}`; + retireWorker(slot, report); } }; worker.addEventListener("message", onFirst); - worker.onerror = (e) => reject(new Error(`occ_service worker: ${e.message}`)); + slot.removeBootListener = () => + worker.removeEventListener("message", onFirst); }); + if (slot.failed || workerSlot !== slot) { + throw new Error("occ_service worker retired during boot"); + } log("[occ] occ_service ready"); - return worker; - })().catch((e) => { - workerP = null; // a failed boot must stay retryable + return slot; + })(); + + slot.ready = Promise.race([boot, bootDeadline]).catch((e) => { + // Do not let a late failure from an old generation clear a replacement + // which a re-entrant caller has already started. + retireWorker(slot, `occ_service unavailable: ${String(e)}`); throw e; }); } - return workerP; + return workerSlot.ready; }; - const post = (worker: Worker, req: OccRequest): Promise => { + const post = (slot: WorkerSlot, req: OccRequest): Promise => { + const worker = slot.worker; + if (!worker || slot.failed || workerSlot !== slot) { + return Promise.resolve({ ok: false, report: "occ_service worker is unavailable" }); + } const id = nextId++; const transfer: Transferable[] = req.kind === "export" ? [req.board.buffer, ...(req.models ?? []).map((m) => m.bytes.buffer)] : [req.bytes.buffer]; return new Promise((resolve) => { - pending.set(id, resolve); - worker.postMessage({ id, req }, transfer); + const timer = setTimeout(() => { + if (pending.get(id)?.generation !== slot.generation) return; + const report = + `occ_service response timed out after ${responseTimeoutMs} ms`; + log(`[occ] ${report} — resetting service`); + retireWorker(slot, report); + }, responseTimeoutMs); + pending.set(id, { generation: slot.generation, resolve, timer }); + try { + worker.postMessage({ id, req }, transfer); + } catch (error) { + pending.delete(id); + clearTimeout(timer); + resolve({ ok: false, report: `occ_service request failed: ${String(error)}` }); + } }); }; + const prefetchBoardModels = async ( + board: Uint8Array, + ): Promise => { + type Outcome = + | { kind: "ready"; models: BoardModelFile[] } + | { kind: "failed"; error: unknown } + | { kind: "timeout" }; + + // collectBoardModelFiles keeps its own bounded network parallelism and does + // no editor-native work. The controller owns this exact optional + // collection: a timeout stops it from selecting more models and makes its + // already-started source results inert. + const controller = new AbortController(); + const collected: Promise = collectBoardModelFiles( + new TextDecoder().decode(board), + 6, + controller.signal, + ).then( + (models) => ({ kind: "ready", models }), + (error) => ({ kind: "failed", error }), + ); + let timer: ReturnType | undefined; + const deadline = new Promise((resolve) => { + timer = setTimeout( + () => { + // Settle the timeout outcome before abort rejection can enqueue its + // Promise reaction, so logs and public behavior stay deterministic. + resolve({ kind: "timeout" }); + controller.abort( + new DOMException("OCC model prefetch timed out", "TimeoutError"), + ); + }, + modelPrefetchTimeoutMs, + ); + }); + const outcome = await Promise.race([collected, deadline]); + if (timer !== undefined) clearTimeout(timer); + + if (outcome.kind === "ready") return outcome.models; + if (!controller.signal.aborted) { + controller.abort( + new DOMException("OCC model prefetch retired", "AbortError"), + ); + } + if (outcome.kind === "failed") { + log(`[occ] model prefetch failed (exporting without models): ${outcome.error}`); + } else { + log( + `[occ] model prefetch timed out after ${modelPrefetchTimeoutMs} ms ` + + "(exporting without models)", + ); + } + return []; + }; + const request = async (req: OccRequest): Promise => { + let prepared: OccRequest; if (req.kind === "export") { + // Capture the caller-owned request fields before the first await and + // build a private dispatch object. A late optional prefetch can then + // neither mutate the caller's object nor change an already-sent payload. + const board = req.board; + const jobJson = req.jobJson; + const fileName = req.fileName; // Ship the board's lib model bodies with the request: the worker's // EXPORTER_STEP resolves them from its own MEMFS (delivery gap doc: // docs/features/3d-models/0007). Best-effort — an export without // models still succeeds, each miss reported by the exporter. - try { - req.models = await collectBoardModelFiles( - new TextDecoder().decode(req.board), - ); - if (req.models.length) - log(`[occ] shipping ${req.models.length} board model(s) with the export`); - } catch (e) { - log(`[occ] model prefetch failed (exporting without models): ${e}`); - req.models = []; - } + const models = await prefetchBoardModels(board); + if (models.length) + log(`[occ] shipping ${models.length} board model(s) with the export`); + prepared = { kind: "export", board, jobJson, fileName, models }; + } else { + prepared = { kind: "loadModel", bytes: req.bytes, ext: req.ext }; } - let worker: Worker; + let slot: WorkerSlot; try { - worker = await ensureWorker(); + slot = await ensureWorker(); } catch (e) { return { ok: false, report: `occ_service unavailable: ${e}` }; } - const res = await post(worker, req); + const res = await post(slot, prepared); - if (req.kind === "export") { + if (prepared.kind === "export") { // Deliver the export straight to the user; the editor gets status only // (the bytes never enter pcbnew's heap). if (res.ok && res.bytes?.length) { @@ -169,7 +382,7 @@ export function installOccService(log: (msg: string) => void): void { // default filename field is empty in the browser); Chromium mangles a // bare dotfile download to "step.txt", so give it a real stem while // keeping the format extension the user picked. - const raw = req.fileName || res.fileName || ""; + const raw = prepared.fileName || res.fileName || ""; const name = !raw || raw.startsWith(".") ? `export${raw || ".step"}` : raw; downloadBytes(name, res.bytes); log(`[occ] export downloaded: ${name} (${res.bytes.length} bytes)`);