findings(E-1,E-2,E-3,E-4): occ-service generations/watchdogs/messageerror + abortable model prefetch

Adapted from codex/asyncify-execution-owner-core 3753320 (scheduler-free on
that branch already; one comment line re-worded for the JSPI line):

- E-1: generation-slotted WorkerSlot with a 2-min boot watchdog (armed before
  resolveWasmBase, so a hung delivery lookup expires too) and a 30-min
  per-request response watchdog; retireWorker() is the single idempotent
  funnel (fail that generation's pendings, terminate, revoke the worker Blob
  URL, clear timers/listeners).
- E-2: worker.onerror is wired for the worker's whole life and settles every
  in-flight STEP/export request; a synchronous postMessage throw settles its
  request without leaking the pending id; late frames from a retired
  generation are inert.
- E-3: worker.onmessageerror retires the generation like error does.
- E-4: collectBoardModelFiles is a pure source/IDB/network path (no editor
  MEMFS round-trip) taking an AbortSignal checked at every loop head;
  prefetchBoardModels races it against a 30 s deadline — timeout is non-fatal
  (export proceeds without models) and late results are inert.

Tests: occ-service.test.ts (7, ported) — boot/response watchdog expiry,
crash-settles-all, bootError retry, decode-fault retirement (invokes the real
onmessageerror transition, per J-4), hung-prefetch export; models-bridge.test.ts
+3 — abort inertness, zero FS access on the collect path. e2e harness twin
updated to the same generation shape (adds __occServiceTestHooks/failDecode).
False-green audited: 19 cases fail with the fixes reverted.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Istvan Matejcsok 2026-08-19 15:48:14 +02:00
commit 3de00bb15b
5 changed files with 957 additions and 97 deletions

View file

@ -45,58 +45,164 @@ export async function installOccServiceStub(page: Page): Promise<void> {
(window as any).__occExports = [];
let workerP: Promise<Worker> | null = null;
const pending = new Map<number, (res: any) => void>();
let nextId = 1;
interface WorkerSlot {
generation: number;
worker?: Worker;
failed: boolean;
ready: Promise<WorkerSlot>;
/** The exact lifecycle transition used by Worker.onmessageerror. */
failDecode: () => void;
}
const ensureWorker = (): Promise<Worker> => {
if (!workerP) {
workerP = (async () => {
interface PendingRequest {
generation: number;
resolve: (res: any) => void;
}
let nextGeneration = 1;
let workerSlot: WorkerSlot | null = null;
const pending = new Map<number, PendingRequest>();
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<WorkerSlot> => {
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<void>((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<void> {
};
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);
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<void> {
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);
}

View file

@ -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<Uint8Array | null>((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", () => {

View file

@ -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, wrlstep 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<BoardModelFile[]> {
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<string>();
let idx = 0;
const worker = async (): Promise<void> => {
while (idx < refs.length) {
while (true) {
throwIfAborted();
if (!isCurrent() || idx >= refs.length) return;
const ref = refs[idx++]!;
for (const candidate of refCandidates(ref)) {
let body: Uint8Array | null = null;
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 });
body = await source.getModelBody(candidate);
} catch {
// best-effort: a missing body surfaces as the exporter's own
// "Could not add 3D model" report warning, never a failed export
// 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;
}

View file

@ -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<MessageListener>();
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<FakeWorker> {
await vi.waitFor(() => expect(FakeWorker.instances.length).toBeGreaterThan(index));
return FakeWorker.instances[index]!;
}
async function readyRequest(workerIndex: number): Promise<{
worker: FakeWorker;
request: Promise<OccResponse>;
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<string>((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" });
});
});

View file

@ -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<number, (res: OccResponse) => void>();
let workerP: Promise<Worker> | 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<Worker> => {
if (!workerP) {
workerP = (async () => {
interface WorkerSlot {
generation: number;
worker?: Worker;
workerUrl?: string;
failed: boolean;
ready: Promise<WorkerSlot>;
bootTimer?: ReturnType<typeof setTimeout>;
rejectBoot?: (reason?: unknown) => void;
removeBootListener?: () => void;
}
interface PendingRequest {
generation: number;
resolve: (res: OccResponse) => void;
timer: ReturnType<typeof setTimeout>;
}
let nextId = 1;
let nextGeneration = 1;
const pending = new Map<number, PendingRequest>();
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<WorkerSlot> => {
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<never>((_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(
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<void>((resolve, reject) => {
await new Promise<void>((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<OccResponse> => {
const post = (slot: WorkerSlot, req: OccRequest): Promise<OccResponse> => {
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<OccResponse>((resolve) => {
pending.set(id, resolve);
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<BoardModelFile[]> => {
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<Outcome> = collectBoardModelFiles(
new TextDecoder().decode(board),
6,
controller.signal,
).then(
(models) => ({ kind: "ready", models }),
(error) => ({ kind: "failed", error }),
);
let timer: ReturnType<typeof setTimeout> | undefined;
const deadline = new Promise<Outcome>((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<OccResponse> => {
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)`);