From ecd55fef95c024279a7d4cd1f87c6b4f8d4b6200 Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Sun, 4 Oct 2026 00:24:32 +0000 Subject: [PATCH 1/5] feat(node): add WebAssembly client transport Co-Authored-By: jason.han --- .github/workflows/pr.yml | 12 +- .../unreleased/node-wasm-transport.added.md | 1 + client/node/README.md | 52 ++ client/node/package.json | 3 +- client/node/src/browser/index.ts | 2 + client/node/src/browser/wasm-worker.ts | 3 + client/node/src/browser/wasm.ts | 243 ++++++++ client/node/src/core/index.ts | 1 + client/node/src/core/wasm.ts | 464 ++++++++++++++ client/node/src/node/index.ts | 2 + client/node/src/node/wasm-worker.ts | 56 ++ client/node/src/node/wasm.ts | 237 +++++++ client/node/test/support/wasm.ts | 89 +++ client/node/test/wasm.test.ts | 585 ++++++++++++++++++ docs/reference/clients.md | 18 +- docs/reference/node-api.md | 2 + docs/reference/wasm.md | 3 + 17 files changed, 1764 insertions(+), 9 deletions(-) create mode 100644 changes/unreleased/node-wasm-transport.added.md create mode 100644 client/node/src/browser/wasm-worker.ts create mode 100644 client/node/src/browser/wasm.ts create mode 100644 client/node/src/core/wasm.ts create mode 100644 client/node/src/node/wasm-worker.ts create mode 100644 client/node/src/node/wasm.ts create mode 100644 client/node/test/support/wasm.ts create mode 100644 client/node/test/wasm.test.ts diff --git a/.github/workflows/pr.yml b/.github/workflows/pr.yml index d9c147413f..d3af5cb9aa 100644 --- a/.github/workflows/pr.yml +++ b/.github/workflows/pr.yml @@ -726,7 +726,7 @@ jobs: runs-on: ubuntu-latest permissions: contents: read - timeout-minutes: 10 + timeout-minutes: 15 steps: - name: Check out repository uses: actions/checkout@v4 @@ -748,6 +748,13 @@ jobs: BUILD_TIME="$(date -u '+%Y-%m-%d_%H:%M:%S')" \ GO_VERSION="$(go version | awk '{print $3}')" + - name: Build the combined WebAssembly module + run: | + mkdir -p bin/wasm + GOOS=js GOARCH=wasm go build -trimpath -ldflags="-s -w" \ + -o bin/wasm/sysml-wasm.wasm ./cmd/sysml-wasm + cp "$(go env GOROOT)/lib/wasm/wasm_exec.js" bin/wasm/ + - name: Verify binaries run: | ./bin/sysml --version @@ -1109,6 +1116,9 @@ jobs: run: | chmod +x bin/sysml-grpc echo "OPENSYSML_BINARY=$(pwd)/bin/sysml-grpc" >> "$GITHUB_ENV" + echo "OPENSYSML_WASM=$(pwd)/bin/wasm/sysml-wasm.wasm" >> "$GITHUB_ENV" + echo "OPENSYSML_WASM_EXEC=$(pwd)/bin/wasm/wasm_exec.js" >> "$GITHUB_ENV" + echo "OPENSYSML_REQUIRE_WASM=1" >> "$GITHUB_ENV" - name: Install the client's dependencies working-directory: client/node diff --git a/changes/unreleased/node-wasm-transport.added.md b/changes/unreleased/node-wasm-transport.added.md new file mode 100644 index 0000000000..3d9384daeb --- /dev/null +++ b/changes/unreleased/node-wasm-transport.added.md @@ -0,0 +1 @@ +- **The Node and browser clients can connect to the combined `sysml-wasm` module without a service.** Node defaults to a worker thread, with inline execution available, and browsers can run inline or in a supplied worker. diff --git a/client/node/README.md b/client/node/README.md index 32ceb5e3ea..4d3dce99c4 100644 --- a/client/node/README.md +++ b/client/node/README.md @@ -171,6 +171,58 @@ Two limits to plan for rather than discover: `fetch` transport, and asserts the allowed origin is answered on the preflight while another origin is not. +## WebAssembly, without a service + +`connectWasm()` uses the same `Connection`, `Model`, values and errors over the +combined `sysml-wasm` module. It serves `ParseFile`, `ParseSources`, +`GetDiagnostics`, `GetSymbol`, `Evaluate`, `Instantiate`, `ExecuteAction`, +`ExecuteState` and `GetServerInfo`. Those are the module's complete RPC surface; +other capability-gated operations fail with `MissingCapabilityError`, and a +direct unsupported RPC fails with `UNIMPLEMENTED`. +The adapter uses JSON encoding; requesting protobuf encoding is refused. + +In Node, `connectWasm()` runs a worker thread by default: + +```ts +import { connectWasm } from "@openmbee/opensysml"; + +await using connection = await connectWasm({ + wasm: "./sysml-wasm.wasm", + wasmExec: "/path/to/the/matching/wasm_exec.js", +}); +const model = await connection.loads("package Demo { part def Car; }"); +``` + +The WASM module and `wasm_exec.js` must come from compatible Go toolchains. +The worker remains referenced while the connection is open, so `close()` or +`await using` ends it. `thread: "inline"` runs Go on the calling thread instead; +it is useful when a worker is unavailable, but blocks that thread during a call. +Closing an inline connection disables its client surface; Go has no exit hook to +stop the running module. +A worker-mode deadline rejects the waiting call without interrupting Go, so +later worker calls queue behind work that outlived its deadline. Inline calls +run synchronously and cannot be interrupted while they block the JavaScript +thread. + +In a browser, omit `worker` to run inline, or provide a module worker serving the +package's `browser/wasm-worker` entry point: + +```ts +import { connectWasm } from "@openmbee/opensysml/browser"; + +const worker = new Worker("/assets/opensysml-wasm-worker.js", { type: "module" }); +await using connection = await connectWasm({ + wasm: new URL("./sysml-wasm.wasm", import.meta.url), + wasmExec: new URL("./wasm_exec.js", import.meta.url), + worker, +}); +``` + +The browser worker module can be bundled from +`@openmbee/opensysml/browser/wasm-worker`. The combined module measures about +7.8 MB gzipped and 5.5 MB with Brotli. A package containing the matching WASM +and Go runtime artifacts will be published separately in a future release. + ## Protobuf, not JSON Bodies are protobuf by default. JSON is available (`connect({ encoding: "json" })`) diff --git a/client/node/package.json b/client/node/package.json index d5a542face..a882dc8b37 100644 --- a/client/node/package.json +++ b/client/node/package.json @@ -29,7 +29,8 @@ "./browser": { "types": "./dist/browser/index.d.ts", "default": "./dist/browser/index.js" - } + }, + "./browser/wasm-worker": "./dist/browser/wasm-worker.js" }, "bin": { "opensysml-generate": "dist/node/generate-cli.js" diff --git a/client/node/src/browser/index.ts b/client/node/src/browser/index.ts index 60f607689b..b5afe9d4aa 100644 --- a/client/node/src/browser/index.ts +++ b/client/node/src/browser/index.ts @@ -8,6 +8,8 @@ import { OpenSysMLError } from "../core/errors.js"; import { baseUrl, encodingOf, interceptors, timeoutOf } from "../core/transport.js"; export * from "../core/index.js"; +export { connectWasm } from "./wasm.js"; +export type { BrowserWasmConnectOptions } from "./wasm.js"; /** How a browser connects: the address is required, because nothing can be started. */ export interface BrowserConnectOptions extends TransportOptions { diff --git a/client/node/src/browser/wasm-worker.ts b/client/node/src/browser/wasm-worker.ts new file mode 100644 index 0000000000..0693cc7734 --- /dev/null +++ b/client/node/src/browser/wasm-worker.ts @@ -0,0 +1,3 @@ +import { serveWasmPort } from "../core/wasm.js"; + +serveWasmPort(globalThis); diff --git a/client/node/src/browser/wasm.ts b/client/node/src/browser/wasm.ts new file mode 100644 index 0000000000..6fb7cae9e8 --- /dev/null +++ b/client/node/src/browser/wasm.ts @@ -0,0 +1,243 @@ +import { Code, ConnectError } from "@connectrpc/connect"; +import { Connection } from "../core/connection.js"; +import type { TransportOptions } from "../core/connection.js"; +import { + connectWasmHost, + instantiateInline, + loadGoConstructor, + type WasmHost, + type WasmWorkerRequest, + type WasmWorkerResponse, + type WorkerWasmSource, +} from "../core/wasm.js"; + +/** Options for connecting to a Go WebAssembly module in a browser. */ +export interface BrowserWasmConnectOptions extends TransportOptions { + /** Module URL, response, bytes, or a compiled WebAssembly module. */ + wasm: string | URL | Response | ArrayBuffer | Uint8Array | WebAssembly.Module; + /** wasm_exec.js from the Go toolchain that built the module. */ + wasmExec?: string | URL; + /** A worker running the package's browser WASM worker entry point. */ + worker?: Worker; + version?: string; + requireCapabilities?: readonly string[]; +} + +/** + * Connects to sysml-wasm inline, or in the provided worker. + */ +export async function connectWasm( + options: BrowserWasmConnectOptions, +): Promise { + const host = + options.worker === undefined + ? await startInline(options.wasm, options.wasmExec) + : await startWorker(options.worker, options.wasm, options.wasmExec); + return connectWasmHost(host, options); +} + +async function startInline( + wasm: BrowserWasmConnectOptions["wasm"], + wasmExec: string | URL | undefined, +): Promise { + const module = await browserWasmSource(wasm); + const Go = await loadGoConstructor(wasmExec === undefined ? undefined : resourceUrl(wasmExec)); + return instantiateInline(module, Go); +} + +async function startWorker( + worker: Worker, + wasm: BrowserWasmConnectOptions["wasm"], + wasmExec: string | URL | undefined, +): Promise { + const host = new BrowserWasmWorkerHost(worker); + try { + const { source, transfer } = await workerSource(wasm); + return await host.start( + { + type: "init", + wasm: source, + ...(wasmExec === undefined ? {} : { wasmExec: resourceUrl(wasmExec) }), + }, + transfer, + ); + } catch (error) { + await host.close(); + throw error; + } +} + +class BrowserWasmWorkerHost implements WasmHost { + version = ""; + + private readonly pending = new Map< + number, + { resolve: (envelope: string) => void; reject: (error: unknown) => void } + >(); + private readonly ready: Promise; + private resolveReady!: () => void; + private rejectReady!: (error: unknown) => void; + private nextId = 1; + private failure: ConnectError | undefined; + private closed = false; + + constructor(private readonly worker: Worker) { + this.ready = new Promise((resolve, reject) => { + this.resolveReady = resolve; + this.rejectReady = reject; + }); + worker.addEventListener("message", (event: MessageEvent) => { + this.receive(event.data); + }); + worker.addEventListener("error", (event: ErrorEvent) => { + this.fail(new Error(event.message || "the sysml-wasm worker failed")); + }); + worker.addEventListener("messageerror", () => { + this.fail(new Error("the sysml-wasm worker sent an unreadable message")); + }); + } + + async start( + request: Extract, + transfer: Transferable[], + ): Promise { + this.worker.postMessage(request, transfer); + await this.ready; + return this; + } + + async call(method: string, params: string): Promise { + this.throwIfFailed(); + if (this.closed) { + throw new ConnectError("the sysml-wasm worker is closed", Code.Unavailable); + } + await this.ready; + this.throwIfFailed(); + const id = this.nextId++; + return new Promise((resolve, reject) => { + this.pending.set(id, { resolve, reject }); + try { + this.worker.postMessage({ type: "call", id, method, params }); + } catch (error) { + this.pending.delete(id); + reject(ConnectError.from(error, Code.Unavailable)); + } + }); + } + + close(): Promise { + if (this.closed) { + return Promise.resolve(); + } + this.closed = true; + try { + this.worker.postMessage({ type: "close" }); + } catch { + // A worker that has already failed has no message loop to close. + } + const closed = new ConnectError("the sysml-wasm worker was closed", Code.Unavailable); + for (const pending of this.pending.values()) { + pending.reject(closed); + } + this.pending.clear(); + this.worker.terminate(); + return Promise.resolve(); + } + + private receive(message: WasmWorkerResponse): void { + if (message.type === "ready") { + this.version = message.version; + this.resolveReady(); + return; + } + if (message.type === "failed") { + this.fail(new Error(message.message)); + return; + } + const pending = this.pending.get(message.id); + if (pending !== undefined) { + this.pending.delete(message.id); + pending.resolve(message.envelope); + } + } + + private fail(reason: unknown): void { + if (this.failure !== undefined || this.closed) { + return; + } + this.failure = ConnectError.from(reason, Code.Unavailable); + this.rejectReady(this.failure); + for (const pending of this.pending.values()) { + pending.reject(this.failure); + } + this.pending.clear(); + } + + private throwIfFailed(): void { + const failure = this.failure; + if (failure !== undefined) { + throw failure; + } + } +} + +async function workerSource( + wasm: BrowserWasmConnectOptions["wasm"], +): Promise<{ source: WorkerWasmSource; transfer: Transferable[] }> { + if (wasm instanceof WebAssembly.Module) { + return { source: wasm, transfer: [] }; + } + if (wasm instanceof Response) { + if (wasm.bodyUsed) { + throw new Error("the WebAssembly Response body has already been read"); + } + const bytes = await wasm.arrayBuffer(); + return { source: bytes, transfer: [bytes] }; + } + if (wasm instanceof ArrayBuffer) { + const bytes = wasm.slice(0); + return { source: bytes, transfer: [bytes] }; + } + if (wasm instanceof Uint8Array) { + const bytes = Uint8Array.from(wasm).buffer; + return { source: bytes, transfer: [bytes] }; + } + const url = wasm instanceof URL ? wasm.href : resourceUrl(wasm); + return { source: url, transfer: [] }; +} + +async function browserWasmSource( + wasm: BrowserWasmConnectOptions["wasm"], +): Promise { + if ( + wasm instanceof WebAssembly.Module || + wasm instanceof ArrayBuffer || + wasm instanceof Uint8Array + ) { + return wasm; + } + const response = + wasm instanceof Response ? wasm : await fetch(wasm instanceof URL ? wasm : resourceUrl(wasm)); + if (!response.ok) { + throw new Error(`could not fetch the WebAssembly module: ${response.status}`); + } + if (response.bodyUsed) { + throw new Error("the WebAssembly Response body has already been read"); + } + if (typeof WebAssembly.compileStreaming === "function") { + try { + return await WebAssembly.compileStreaming(response.clone()); + } catch { + // Some servers omit the application/wasm content type required by streaming. + } + } + return response.arrayBuffer(); +} + +function resourceUrl(value: string | URL): string { + if (value instanceof URL) { + return value.href; + } + const base = typeof document === "undefined" ? import.meta.url : document.baseURI; + return new URL(value, base).href; +} diff --git a/client/node/src/core/index.ts b/client/node/src/core/index.ts index 8c2cef3f7e..22a8b05839 100644 --- a/client/node/src/core/index.ts +++ b/client/node/src/core/index.ts @@ -7,6 +7,7 @@ export type { ResponseTap, TransportOptions, } from "./connection.js"; +export type { WasmHost } from "./wasm.js"; export { Instance, InstanceTree, Model, ModelSymbol, decodeDiagnostic } from "./model.js"; export type { AttributeFacts, diff --git a/client/node/src/core/wasm.ts b/client/node/src/core/wasm.ts new file mode 100644 index 0000000000..77521f37d0 --- /dev/null +++ b/client/node/src/core/wasm.ts @@ -0,0 +1,464 @@ +import { + fromJson, + toJson, + type DescMessage, + type DescMethodUnary, + type JsonValue, +} from "@bufbuild/protobuf"; +import { + Code, + ConnectError, + createContextValues, + type StreamResponse, + type Transport, + type UnaryResponse, +} from "@connectrpc/connect"; +import { getAbortSignalReason, runUnaryCall } from "@connectrpc/connect/protocol"; +import { Connection } from "./connection.js"; +import type { TransportOptions } from "./connection.js"; +import { ClosedConnectionError, OpenSysMLError } from "./errors.js"; +import { encodingOf, interceptors, timeoutOf } from "./transport.js"; + +/** A host surface implemented by an inline Go runtime or a worker. */ +export interface WasmHost { + readonly version: string; + call(method: string, params: string): Promise; + close(): Promise; +} + +/** The Go runtime interface exposed by wasm_exec.js. */ +export interface GoRuntime { + readonly importObject: WebAssembly.Imports; + run(instance: WebAssembly.Instance): Promise; +} + +/** A constructor exported globally by wasm_exec.js. */ +export type GoConstructor = new () => GoRuntime; + +export type WorkerWasmSource = string | ArrayBuffer | Uint8Array | WebAssembly.Module; + +export type WasmWorkerRequest = + | { type: "init"; wasm: WorkerWasmSource; wasmExec?: string } + | { type: "call"; id: number; method: string; params: string } + | { type: "close" }; + +export type WasmWorkerResponse = + | { type: "ready"; version: string } + | { type: "failed"; message: string } + | { type: "answer"; id: number; envelope: string }; + +/** The message port surface shared by browser and Node workers. */ +export interface WasmPortLike { + postMessage(message: WasmWorkerResponse): void; + addEventListener( + type: "message", + listener: (event: MessageEvent) => void, + ): void; + start?(): void; +} + +/** Optional platform loaders for worker runtimes that need local file access. */ +export interface WasmPortLoaders { + loadWasm?(source: WorkerWasmSource): Promise; + loadGo?(wasmExec: string | undefined): Promise; +} + +/** Connection settings shared by the Node and browser WASM clients. */ +export interface WasmConnectionOptions extends TransportOptions { + version?: string; + requireCapabilities?: readonly string[]; +} + +interface GlobalWasmSurface { + readonly version: string; + call(method: string, params: string): string; +} + +let inlineStartup = Promise.resolve(); + +/** + * Adapts the Go host global to the asynchronous host interface. + */ +export function wasmHostFromGlobal(surface?: GlobalWasmSurface): WasmHost { + const global = globalThis as typeof globalThis & { sysmlWasm?: GlobalWasmSurface }; + const hostSurface = surface ?? global.sysmlWasm; + if ( + hostSurface === undefined || + typeof hostSurface.version !== "string" || + typeof hostSurface.call !== "function" + ) { + throw new OpenSysMLError("the Go runtime did not expose globalThis.sysmlWasm"); + } + let closed = false; + return { + version: hostSurface.version, + call(method, params) { + if (closed) { + return Promise.reject(new ClosedConnectionError()); + } + try { + return Promise.resolve(hostSurface.call(method, params)); + } catch (cause) { + return Promise.resolve(errorEnvelope(cause)); + } + }, + close() { + closed = true; + return Promise.resolve(); + }, + }; +} + +/** + * Starts a Go WebAssembly module in this realm and captures its host global. + */ +export async function instantiateInline( + module: BufferSource | WebAssembly.Module, + Go: GoConstructor, +): Promise { + const starting = inlineStartup.then(() => instantiateInlineNow(module, Go)); + inlineStartup = starting.then( + () => undefined, + () => undefined, + ); + return await starting; +} + +async function instantiateInlineNow( + module: BufferSource | WebAssembly.Module, + Go: GoConstructor, +): Promise { + const runtime = globalThis as typeof globalThis & { sysmlWasm?: GlobalWasmSurface }; + delete runtime.sysmlWasm; + + const go = new Go(); + const instance = + module instanceof WebAssembly.Module + ? await WebAssembly.instantiate(module, go.importObject) + : (await WebAssembly.instantiate(module, go.importObject)).instance; + let startError: unknown; + void go.run(instance).catch((error: unknown) => { + startError = error; + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + + const surface = globalWasmSurface(); + delete runtime.sysmlWasm; + if (surface === undefined) { + throw new OpenSysMLError( + "the Go runtime did not expose globalThis.sysmlWasm", + startError === undefined ? undefined : { cause: startError }, + ); + } + return wasmHostFromGlobal(surface); +} + +/** + * Creates a Connect transport for the JSON envelope served by sysml-wasm. + */ +export function createWasmTransport(host: WasmHost, options: TransportOptions): Transport { + const callUnary = ( + method: DescMethodUnary, + signal: AbortSignal | undefined, + timeoutMs: number | undefined, + header: HeadersInit | undefined, + input: Parameters[4], + contextValues: Parameters[5], + ): Promise> => + runUnaryCall({ + req: { + stream: false, + service: method.parent, + method, + requestMethod: "POST", + url: `wasm://sysml-wasm/${method.parent.typeName}/${method.name}`, + header: new Headers(header), + contextValues: contextValues ?? createContextValues(), + message: input as never, + }, + ...(signal === undefined ? {} : { signal }), + ...(timeoutMs === undefined ? {} : { timeoutMs }), + interceptors: interceptors(options), + next: async (request) => { + if (request.signal.aborted) { + throw abortError(request.signal); + } + const params = JSON.stringify(toJson(method.input, request.message)); + const answer = host.call(method.name, params); + const envelope = await withAbort(answer, request.signal); + const body = readEnvelope(envelope); + return { + stream: false, + service: method.parent, + method, + header: new Headers(), + message: fromJson(method.output, body as JsonValue, { + ignoreUnknownFields: true, + }), + trailer: new Headers(), + }; + }, + }); + + return { + unary: callUnary, + stream(): Promise> { + return Promise.reject( + new ConnectError("streaming methods are not supported by sysml-wasm", Code.Unimplemented), + ); + }, + }; +} + +/** + * Opens a Connection over a WASM host, negotiating its version and capabilities. + */ +export async function connectWasmHost( + host: WasmHost, + options: WasmConnectionOptions = {}, +): Promise { + let timeoutMs: number | undefined; + try { + if (options.encoding !== undefined && encodingOf(options) !== "json") { + throw new OpenSysMLError("sysml-wasm supports JSON encoding only"); + } + timeoutMs = timeoutOf(options); + } catch (error) { + await host.close().catch(() => undefined); + throw error; + } + const required = + options.version === undefined || options.version === "" || options.version === "latest" + ? undefined + : options.version; + const origin = `sysml-wasm ${host.version}`; + return Connection.open({ + transport: createWasmTransport(host, options), + backend: { + origin, + release: () => host.close(), + warn: (message) => { + console.warn(message); + }, + }, + encoding: "json", + timeoutMs, + ...(required === undefined ? {} : { requiredVersion: required }), + ...(options.requireCapabilities === undefined + ? {} + : { requiredCapabilities: options.requireCapabilities }), + stale: { + address: origin, + remedy: required === undefined + ? "use a module reporting the expected version, or omit version" + : `use a module reporting ${required}, or omit version`, + }, + }); +} + +/** + * Serves the shared init/call protocol on a Node or browser worker port. + */ +export function serveWasmPort(port: WasmPortLike, loaders: WasmPortLoaders = {}): void { + let host: WasmHost | undefined; + let initializing = false; + let closed = false; + let calls = Promise.resolve(); + + const onMessage = (event: MessageEvent): void => { + const request = event.data; + if (request.type === "init") { + if (initializing || host !== undefined || closed) { + port.postMessage({ type: "failed", message: "the WASM worker is already initialized" }); + return; + } + initializing = true; + void initialize(request.wasm, request.wasmExec); + return; + } + if (request.type === "close") { + closed = true; + void host?.close(); + return; + } + calls = calls.then(async () => { + if (closed || host === undefined) { + port.postMessage({ + type: "answer", + id: request.id, + envelope: errorEnvelope(new ConnectError("the WASM worker is unavailable", Code.Unavailable)), + }); + return; + } + try { + port.postMessage({ + type: "answer", + id: request.id, + envelope: await host.call(request.method, request.params), + }); + } catch (cause) { + port.postMessage({ type: "answer", id: request.id, envelope: errorEnvelope(cause) }); + } + }); + }; + + const initialize = async ( + source: WorkerWasmSource, + wasmExec: string | undefined, + ): Promise => { + try { + const module = loaders.loadWasm + ? await loaders.loadWasm(source) + : await loadWasmSource(source); + const Go = loaders.loadGo + ? await loaders.loadGo(wasmExec) + : await loadGoConstructor(wasmExec); + host = await instantiateInline(module, Go); + port.postMessage({ type: "ready", version: host.version }); + } catch (cause) { + port.postMessage({ type: "failed", message: errorMessage(cause) }); + } + }; + + port.addEventListener("message", onMessage); + port.start?.(); +} + +/** Loads the Go constructor installed by wasm_exec.js, importing it when needed. */ +export async function loadGoConstructor( + wasmExec?: string, + forceImport = false, +): Promise { + const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; + if (wasmExec !== undefined && (forceImport || runtime.Go === undefined)) { + try { + await import(wasmExec); + } catch (cause) { + const importScripts = ( + globalThis as typeof globalThis & { importScripts?: (...urls: string[]) => void } + ).importScripts; + if (importScripts === undefined) { + throw new OpenSysMLError(`could not load wasm_exec.js from ${wasmExec}`, { cause }); + } + importScripts(wasmExec); + } + } + if (runtime.Go === undefined) { + throw new OpenSysMLError("Go is unavailable; supply the matching wasm_exec.js"); + } + return runtime.Go; +} + +function readEnvelope(envelope: string): unknown { + let parsed: unknown; + try { + parsed = JSON.parse(envelope) as unknown; + } catch (cause) { + throw ConnectError.from(cause, Code.Internal); + } + if (parsed === null || typeof parsed !== "object") { + throw new ConnectError("sysml-wasm returned an invalid JSON envelope", Code.Internal); + } + const response = parsed as { + result?: unknown; + error?: { code?: unknown; message?: unknown } | null; + }; + if (response.error !== undefined && response.error !== null) { + const status = + typeof response.error.code === "number" && + Object.values(Code).includes(response.error.code) + ? response.error.code + : Code.Unknown; + throw new ConnectError( + typeof response.error.message === "string" + ? response.error.message + : "sysml-wasm returned an error", + status, + ); + } + if (!Object.hasOwn(response, "result")) { + throw new ConnectError("sysml-wasm returned an envelope without a result", Code.Internal); + } + return response.result; +} + +function withAbort(answer: Promise, signal: AbortSignal): Promise { + if (signal.aborted) { + return Promise.reject(abortError(signal)); + } + return new Promise((resolve, reject) => { + let settled = false; + const cleanup = (): void => { + signal.removeEventListener("abort", onAbort); + }; + const finish = (callback: (value: T) => void, value: T): void => { + if (!settled) { + settled = true; + cleanup(); + callback(value); + } + }; + const fail = (error: unknown): void => { + if (!settled) { + settled = true; + cleanup(); + reject(ConnectError.from(error, Code.Internal)); + } + }; + const onAbort = (): void => { + fail(abortError(signal)); + }; + signal.addEventListener("abort", onAbort, { once: true }); + answer.then( + (value) => { + finish(resolve, value); + }, + fail, + ); + }); +} + +function abortError(signal: AbortSignal): ConnectError { + return ConnectError.from(getAbortSignalReason(signal), Code.Canceled); +} + +function errorEnvelope(error: unknown): string { + const connectError = ConnectError.from(error, Code.Internal); + return JSON.stringify({ + jsonrpc: "2.0", + id: null, + error: { code: connectError.code, message: connectError.rawMessage }, + }); +} + +async function loadWasmSource( + source: WorkerWasmSource, +): Promise { + if ( + source instanceof WebAssembly.Module || + source instanceof ArrayBuffer || + ArrayBuffer.isView(source) + ) { + return source; + } + const response = await fetch(source); + if (!response.ok) { + throw new OpenSysMLError(`could not fetch the WebAssembly module: ${response.status}`); + } + if (typeof WebAssembly.compileStreaming === "function") { + try { + return await WebAssembly.compileStreaming(response.clone()); + } catch { + // Some servers omit the application/wasm content type required by streaming. + } + } + return response.arrayBuffer(); +} + +function errorMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} + +function globalWasmSurface(): GlobalWasmSurface | undefined { + return (globalThis as typeof globalThis & { sysmlWasm?: GlobalWasmSurface }).sysmlWasm; +} diff --git a/client/node/src/node/index.ts b/client/node/src/node/index.ts index 3a0b303779..c5d5959df5 100644 --- a/client/node/src/node/index.ts +++ b/client/node/src/node/index.ts @@ -89,6 +89,8 @@ export { verifyManifest, } from "./signing.js"; export { PrivateService, currentPrivateService } from "./service.js"; +export { connectWasm } from "./wasm.js"; +export type { WasmConnectOptions } from "./wasm.js"; /** Names a service to connect to instead of starting one. */ export const SERVICE_ENV = "OPENSYSML_SERVICE"; diff --git a/client/node/src/node/wasm-worker.ts b/client/node/src/node/wasm-worker.ts new file mode 100644 index 0000000000..4a68ed8761 --- /dev/null +++ b/client/node/src/node/wasm-worker.ts @@ -0,0 +1,56 @@ +import { readFile } from "node:fs/promises"; +import { resolve } from "node:path"; +import { pathToFileURL } from "node:url"; +import { parentPort } from "node:worker_threads"; +import { + loadGoConstructor, + serveWasmPort, + type WasmPortLike, + type WorkerWasmSource, +} from "../core/wasm.js"; + +if (parentPort === null) { + throw new Error("the sysml-wasm worker must run in a worker thread"); +} + +serveWasmPort(parentPort as unknown as WasmPortLike, { + loadWasm: loadNodeWasm, + loadGo: (wasmExec) => loadGoConstructor(wasmExec === undefined ? undefined : moduleSpecifier(wasmExec)), +}); + +async function loadNodeWasm(source: WorkerWasmSource): Promise { + if ( + source instanceof WebAssembly.Module || + source instanceof ArrayBuffer || + ArrayBuffer.isView(source) + ) { + return source; + } + const url = parseUrl(source); + if (url?.protocol === "file:") { + return readFile(url); + } + if (url !== undefined) { + const response = await fetch(url); + if (!response.ok) { + throw new Error(`could not fetch the WebAssembly module: ${response.status}`); + } + return response.arrayBuffer(); + } + return readFile(resolve(source)); +} + +function moduleSpecifier(source: string): string { + return parseUrl(source)?.href ?? pathToFileURL(resolve(source)).href; +} + +function parseUrl(value: string): URL | undefined { + if (/^[A-Za-z]:[\\/]/.test(value)) { + return undefined; + } + try { + return new URL(value); + } catch { + return undefined; + } +} diff --git a/client/node/src/node/wasm.ts b/client/node/src/node/wasm.ts new file mode 100644 index 0000000000..9ce87541ae --- /dev/null +++ b/client/node/src/node/wasm.ts @@ -0,0 +1,237 @@ +import { readFile } from "node:fs/promises"; +import { isAbsolute, resolve } from "node:path"; +import { pathToFileURL } from "node:url"; +import { Worker } from "node:worker_threads"; +import { Code, ConnectError } from "@connectrpc/connect"; +import { Connection } from "../core/connection.js"; +import type { TransportOptions } from "../core/connection.js"; +import { + connectWasmHost, + instantiateInline, + loadGoConstructor, + type WasmHost, + type WasmWorkerRequest, + type WasmWorkerResponse, + type WorkerWasmSource, +} from "../core/wasm.js"; + +/** Options for connecting to a Go WebAssembly module from Node. */ +export interface WasmConnectOptions extends TransportOptions { + /** Path, file URL, HTTP URL, or bytes of sysml-wasm.wasm. */ + wasm: string | URL | Uint8Array; + /** wasm_exec.js from the Go toolchain that built the module. */ + wasmExec: string | URL; + /** Run the module in a worker thread (default) or this thread. */ + thread?: "worker" | "inline"; + version?: string; + requireCapabilities?: readonly string[]; +} + +/** + * Connects to sysml-wasm, in a worker by default or inline when requested. + */ +export async function connectWasm(options: WasmConnectOptions): Promise { + const wasmExec = moduleSpecifier(options.wasmExec); + let host: WasmHost; + if (options.thread === "inline") { + const module = await loadNodeWasm(options.wasm); + const Go = await loadGoConstructor(wasmExec, true); + host = await instantiateInline(module, Go); + } else { + host = await startWorker(options.wasm, wasmExec); + } + return connectWasmHost(host, options); +} + +async function startWorker(wasm: WasmConnectOptions["wasm"], wasmExec: string): Promise { + const worker = new Worker(new URL("./wasm-worker.js", import.meta.url)); + const host = new NodeWasmWorkerHost(worker); + const { source, transfer } = workerSource(wasm); + try { + return await host.start( + { type: "init", wasm: source, wasmExec }, + transfer, + ); + } catch (error) { + await host.close(); + throw error; + } +} + +class NodeWasmWorkerHost implements WasmHost { + version = ""; + + private readonly pending = new Map< + number, + { resolve: (envelope: string) => void; reject: (error: unknown) => void } + >(); + private readonly ready: Promise; + private resolveReady!: () => void; + private rejectReady!: (error: unknown) => void; + private nextId = 1; + private failure: ConnectError | undefined; + private closed = false; + + constructor(private readonly worker: Worker) { + this.ready = new Promise((resolve, reject) => { + this.resolveReady = resolve; + this.rejectReady = reject; + }); + worker.on("message", (message: WasmWorkerResponse) => { + this.receive(message); + }); + worker.on("error", (error: Error) => { + this.fail(error); + }); + worker.on("exit", (code: number) => { + if (!this.closed) { + this.fail(new Error(`the sysml-wasm worker exited with code ${code}`)); + } + }); + } + + async start( + request: Extract, + transfer: readonly ArrayBuffer[], + ): Promise { + this.worker.postMessage(request, transfer); + await this.ready; + return this; + } + + async call(method: string, params: string): Promise { + this.throwIfFailed(); + if (this.closed) { + throw new ConnectError("the sysml-wasm worker is closed", Code.Unavailable); + } + await this.ready; + this.throwIfFailed(); + const id = this.nextId++; + return new Promise((resolve, reject) => { + this.pending.set(id, { resolve, reject }); + try { + this.worker.postMessage({ type: "call", id, method, params }); + } catch (error) { + this.pending.delete(id); + reject(ConnectError.from(error, Code.Unavailable)); + } + }); + } + + async close(): Promise { + if (this.closed) { + return; + } + this.closed = true; + try { + this.worker.postMessage({ type: "close" }); + } catch { + // A worker that has already failed has no message loop to close. + } + const closed = new ConnectError("the sysml-wasm worker was closed", Code.Unavailable); + for (const pending of this.pending.values()) { + pending.reject(closed); + } + this.pending.clear(); + await this.worker.terminate(); + } + + private receive(message: WasmWorkerResponse): void { + if (message.type === "ready") { + this.version = message.version; + this.resolveReady(); + return; + } + if (message.type === "failed") { + this.fail(new Error(message.message)); + return; + } + const pending = this.pending.get(message.id); + if (pending !== undefined) { + this.pending.delete(message.id); + pending.resolve(message.envelope); + } + } + + private fail(reason: unknown): void { + if (this.failure !== undefined || this.closed) { + return; + } + this.failure = ConnectError.from(reason, Code.Unavailable); + this.rejectReady(this.failure); + for (const pending of this.pending.values()) { + pending.reject(this.failure); + } + this.pending.clear(); + } + + private throwIfFailed(): void { + const failure = this.failure; + if (failure !== undefined) { + throw failure; + } + } +} + +function workerSource( + wasm: WasmConnectOptions["wasm"], +): { source: WorkerWasmSource; transfer: readonly ArrayBuffer[] } { + if (wasm instanceof Uint8Array) { + const bytes = Uint8Array.from(wasm); + return { source: bytes.buffer, transfer: [bytes.buffer] }; + } + return { source: wasm instanceof URL ? wasm.href : wasm, transfer: [] }; +} + +async function loadNodeWasm(wasm: WasmConnectOptions["wasm"]): Promise { + if (wasm instanceof Uint8Array) { + return wasm; + } + if (wasm instanceof URL) { + if (wasm.protocol === "file:") { + return readFile(wasm); + } + const response = await fetch(wasm); + if (!response.ok) { + throw new Error(`could not fetch the WebAssembly module: ${response.status}`); + } + return response.arrayBuffer(); + } + const url = parseUrl(wasm); + if (url?.protocol === "file:") { + return readFile(url); + } + if (url !== undefined) { + const response = await fetch(url); + if (!response.ok) { + throw new Error(`could not fetch the WebAssembly module: ${response.status}`); + } + return response.arrayBuffer(); + } + return readFile(wasm); +} + +function moduleSpecifier(source: string | URL): string { + if (source instanceof URL) { + return source.href; + } + if (!/^[A-Za-z]:[\\/]/.test(source)) { + try { + return new URL(source).href; + } catch { + // Resolve filesystem paths relative to the current working directory. + } + } + return pathToFileURL(isAbsolute(source) ? source : resolve(source)).href; +} + +function parseUrl(source: string): URL | undefined { + if (/^[A-Za-z]:[\\/]/.test(source)) { + return undefined; + } + try { + return new URL(source); + } catch { + return undefined; + } +} diff --git a/client/node/test/support/wasm.ts b/client/node/test/support/wasm.ts new file mode 100644 index 0000000000..ab52769c5d --- /dev/null +++ b/client/node/test/support/wasm.ts @@ -0,0 +1,89 @@ +import { execFileSync } from "node:child_process"; +import { + existsSync, + mkdirSync, + renameSync, + unlinkSync, +} from "node:fs"; +import { join } from "node:path"; +import { packageRoot, repoRoot } from "./service.js"; + +export interface WasmArtifacts { + wasm: string; + wasmExec: string; +} + +let artifacts: Promise | undefined; + +/** Uses supplied WebAssembly artifacts or builds the combined module once. */ +export function wasmArtifacts(): Promise { + artifacts ??= Promise.resolve().then(resolveArtifacts); + return artifacts; +} + +function resolveArtifacts(): WasmArtifacts | undefined { + const wasm = process.env.OPENSYSML_WASM; + const wasmExec = process.env.OPENSYSML_WASM_EXEC; + if (wasm !== undefined || wasmExec !== undefined) { + if (wasm === undefined || wasmExec === undefined) { + throw new Error("OPENSYSML_WASM and OPENSYSML_WASM_EXEC must be set together"); + } + if (!existsSync(wasm) || !existsSync(wasmExec)) { + throw new Error("the configured WebAssembly module or wasm_exec.js does not exist"); + } + return { wasm, wasmExec }; + } + + let goroot: string; + try { + goroot = execFileSync("go", ["env", "GOROOT"], { encoding: "utf8" }).trim(); + } catch (cause) { + if (isMissingGo(cause) && process.env.OPENSYSML_REQUIRE_WASM !== "1") { + return undefined; + } + throw new Error("Go is required to build the sysml-wasm test module", { cause }); + } + + const dir = join(packageRoot, "build", "wasm"); + const output = join(dir, "sysml-wasm.wasm"); + const runtime = join(goroot, "lib", "wasm", "wasm_exec.js"); + if (!existsSync(runtime)) { + throw new Error(`the Go runtime script does not exist: ${runtime}`); + } + if (!existsSync(output)) { + mkdirSync(dir, { recursive: true }); + const staged = `${output}.${String(process.pid)}.stage`; + try { + execFileSync( + "go", + [ + "build", + "-trimpath", + "-ldflags=-s -w", + "-o", + staged, + "./cmd/sysml-wasm", + ], + { + cwd: repoRoot, + env: { ...process.env, GOOS: "js", GOARCH: "wasm" }, + stdio: "inherit", + }, + ); + renameSync(staged, output); + } finally { + if (existsSync(staged)) { + unlinkSync(staged); + } + } + } + return { wasm: output, wasmExec: runtime }; +} + +function isMissingGo(error: unknown): boolean { + return ( + error instanceof Error && + "code" in error && + (error as NodeJS.ErrnoException).code === "ENOENT" + ); +} diff --git a/client/node/test/wasm.test.ts b/client/node/test/wasm.test.ts new file mode 100644 index 0000000000..902d69725c --- /dev/null +++ b/client/node/test/wasm.test.ts @@ -0,0 +1,585 @@ +import assert from "node:assert/strict"; +import { readFile } from "node:fs/promises"; +import { join } from "node:path"; +import { pathToFileURL } from "node:url"; +import { after, test } from "node:test"; +import { Worker as NodeWorker } from "node:worker_threads"; +import { Code, ConnectError, createClient } from "@connectrpc/connect"; +import { + ClosedConnectionError, + MissingCapabilityError, + OpenSysMLError, + SourceDocument, + StaleServiceError, + connect, + connectWasm, +} from "../src/node/index.js"; +import { connectWasm as connectBrowserWasm } from "../src/browser/index.js"; +import { + connectWasmHost, + createWasmTransport, + instantiateInline, + serveWasmPort, + type GoConstructor, + type WasmHost, + type WasmPortLike, + type WasmWorkerRequest, + type WasmWorkerResponse, +} from "../src/core/wasm.js"; +import { SysMLService } from "../src/generated/sysml_pb.js"; +import { wasmArtifacts, type WasmArtifacts } from "./support/wasm.js"; +import { repoRoot, SAMPLE, useServiceBinary } from "./support/service.js"; + +const artifacts = await wasmArtifacts(); +if (artifacts !== undefined) { + useServiceBinary(); +} +const wasmSkip = + artifacts === undefined ? "Go is unavailable; the WASM integration test was skipped" : false; + +after(() => { + Reflect.deleteProperty(globalThis, "sysmlWasm"); +}); + +test("the transport maps JSON envelopes, ignores unknown fields, and observes aborts", async () => { + const unknownFields = clientFor( + fixedHost( + '{"jsonrpc":"2.0","id":null,"result":{"version":"transport-test","capabilities":[],"futureField":true}}', + ), + ); + const response = await unknownFields.getServerInfo({}); + assert.equal(response.version, "transport-test"); + assert.deepEqual(response.capabilities, []); + + const failed = clientFor( + fixedHost( + '{"jsonrpc":"2.0","id":null,"error":{"code":5,"message":"not found"}}', + ), + ); + await assert.rejects( + () => failed.getServerInfo({}), + (error: unknown) => + error instanceof ConnectError && error.code === Code.NotFound, + ); + + const controller = new AbortController(); + let closeCount = 0; + const unanswered = clientFor({ + ...fixedHost(""), + call: () => new Promise(() => {}), + close: () => { + closeCount += 1; + return Promise.resolve(); + }, + }); + const pending = unanswered.getServerInfo({}, { signal: controller.signal }); + controller.abort(new ConnectError("cancelled", Code.Canceled)); + await assert.rejects( + pending, + (error: unknown) => error instanceof ConnectError && error.code === Code.Canceled, + ); + assert.equal(closeCount, 0); + + const deadline = clientFor({ + ...fixedHost(""), + call: () => new Promise(() => {}), + close: () => { + closeCount += 1; + return Promise.resolve(); + }, + }); + await assert.rejects( + () => deadline.getServerInfo({}, { timeoutMs: 20 }), + (error: unknown) => + error instanceof ConnectError && error.code === Code.DeadlineExceeded, + ); + assert.equal(closeCount, 0); +}); + +test("invalid WASM connection settings close the host before handshake", async () => { + for (const options of [{ encoding: "protobuf" as const }, { timeoutMs: 0 }]) { + let closed = false; + const host: WasmHost = { + ...fixedHost(""), + close: () => { + closed = true; + return Promise.resolve(); + }, + }; + await assert.rejects(() => connectWasmHost(host, options), OpenSysMLError); + assert.equal(closed, true); + } +}); + +test("inline hosts capture the Go surface and reject calls after close", async () => { + const host = await instantiateInline(emptyWasm, FakeGo); + assert.equal("sysmlWasm" in globalThis, false); + assert.equal(host.version, "fake-wasm"); + assert.deepEqual(JSON.parse(await host.call("Echo", "{}")), { + jsonrpc: "2.0", + id: null, + result: { method: "Echo", params: {} }, + }); + await host.close(); + await assert.rejects(host.call("Echo", "{}"), ClosedConnectionError); +}); + +test("worker ports initialize, answer calls, and return initialization errors", async () => { + const pair = fakePortPair(); + const received: WasmWorkerResponse[] = []; + pair.client.addEventListener("message", (event) => received.push(event.data)); + serveWasmPort(pair.service, { + loadWasm: () => Promise.resolve(emptyWasm), + loadGo: () => Promise.resolve(FakeGo), + }); + pair.client.postMessage({ type: "init", wasm: emptyWasm }); + const ready = await waitFor(received, (message) => message.type === "ready"); + assert.equal(ready.version, "fake-wasm"); + + pair.client.postMessage({ type: "call", id: 1, method: "Echo", params: "{}" }); + const answer = await waitFor( + received, + (message): message is Extract => + message.type === "answer" && message.id === 1, + ); + assert.deepEqual( + (JSON.parse(answer.envelope) as { result: unknown }).result, + { method: "Echo", params: {} }, + ); + + pair.client.postMessage({ type: "call", id: 2, method: "Fail", params: "{}" }); + const failedCall = await waitFor( + received, + (message): message is Extract => + message.type === "answer" && message.id === 2, + ); + assert.equal( + (JSON.parse(failedCall.envelope) as { error: { code: number } }).error.code, + Code.Internal, + ); + + const broken = fakePortPair(); + const brokenMessages: WasmWorkerResponse[] = []; + broken.client.addEventListener("message", (event) => { + brokenMessages.push(event.data); + }); + serveWasmPort(broken.service, { + loadWasm: () => Promise.reject(new Error("invalid module")), + }); + broken.client.postMessage({ type: "init", wasm: emptyWasm }); + const initError = await waitFor(brokenMessages, (message) => message.type === "failed"); + assert.deepEqual(initError, { type: "failed", message: "invalid module" }); +}); + +test("browser WASM workers connect through the shared request protocol", async () => { + const sources = [ + emptyWasm, + emptyWasm.buffer.slice(0), + new WebAssembly.Module(emptyWasm), + new Response(emptyWasm), + ]; + for (const wasm of sources) { + const worker = new FakeBrowserWorker(); + const connection = await connectBrowserWasm({ + wasm, + worker: worker as unknown as Worker, + }); + assert.equal((await connection.rpc.getServerInfo({})).version, "fake-wasm"); + await connection.close(); + assert.equal(worker.terminated, true); + } + assert.notEqual(emptyWasm.byteLength, 0); +}); + +for (const thread of ["worker", "inline"] as const) { + test(`Node ${thread} WASM connections match the native client`, { skip: wasmSkip }, async () => { + const wasmFiles = requireArtifacts(); + const wasmResponses: string[] = []; + const wasmBytes = + thread === "worker" ? new Uint8Array(await readFile(wasmFiles.wasm)) : undefined; + await using wasm = await connectWasm({ + wasm: wasmBytes ?? wasmFiles.wasm, + wasmExec: wasmFiles.wasmExec, + thread, + onResponse: ({ method }) => wasmResponses.push(method), + }); + if (wasmBytes !== undefined) { + assert.notEqual(wasmBytes.byteLength, 0); + } + await using native = await connect(); + const actualCapabilities = [...wasm.info.capabilities]; + assert.notEqual(wasm.info.version, ""); + assert.deepEqual(actualCapabilities, [ + "type_facts", + "enum_values", + "evaluate_subject", + "symbol_attributes", + "unset_value", + "feature_values", + "inline_language", + "strict_conformance", + "parse_sources", + "complex_values", + "structured_values", + "measurement_refs", + "function_values", + "set_values", + "tensor_values", + "infinity_value", + "diagnostic_codes", + "schedule", + "final_time", + "metaobject_values", + "undetermined_value", + "performer", + "big_int_values", + ]); + + const wasmModel = await wasm.loads(SAMPLE); + const nativeModel = await native.loads(SAMPLE); + assert.equal(wasmModel.hash, nativeModel.hash); + assert.deepEqual(wasmModel.diagnostics, nativeModel.diagnostics); + const wasmCar = await wasmModel.symbol("Sample::Car"); + const nativeCar = await nativeModel.symbol("Sample::Car"); + assert.equal(wasmCar.id, nativeCar.id); + assert.equal(wasmCar.kind, nativeCar.kind); + assert.deepEqual( + (await wasmCar.children()).map((child) => child.name), + (await nativeCar.children()).map((child) => child.name), + ); + assert.deepEqual( + await wasmModel.refreshDiagnostics(), + await nativeModel.refreshDiagnostics(), + ); + assert.deepEqual(await wasmModel.eval("2 + 2"), await nativeModel.eval("2 + 2")); + assert.deepEqual( + await wasmModel.eval("Sample::Car::mass"), + await nativeModel.eval("Sample::Car::mass"), + ); + assert.deepEqual( + await wasmModel.instantiate("Sample::Car"), + await nativeModel.instantiate("Sample::Car"), + ); + + const validationText = await readFile( + join(repoRoot, "tests", "wasm", "testdata", "core-validation.sysml"), + "utf8", + ); + const wasmValidation = await wasm.loads(validationText); + const nativeValidation = await native.loads(validationText); + assert.deepEqual(wasmValidation.diagnostics, nativeValidation.diagnostics); + assert.ok(wasmValidation.diagnostics.some((diagnostic) => diagnostic.severity === "error")); + + const engineModelText = await readFile( + join(repoRoot, "tests", "wasm", "testdata", "engine.sysml"), + "utf8", + ); + const wasmEngineModel = await wasm.loads(engineModelText); + const nativeEngineModel = await native.loads(engineModelText); + for (const expression of ["7", "1.5", "true", "2 ** 70"]) { + assert.deepEqual( + await wasmEngineModel.eval(expression), + await nativeEngineModel.eval(expression), + `${thread} result for ${expression}`, + ); + } + const wasmQuantity = await wasmEngineModel.eval("enginedemo::speed"); + assert.equal(wasmQuantity.kind, "quantity"); + assert.deepEqual( + wasmQuantity, + await nativeEngineModel.eval("enginedemo::speed"), + ); + assert.deepEqual( + await wasmEngineModel.executeAction("enginedemo::Double", { inputs: { x: 6n } }), + await nativeEngineModel.executeAction("enginedemo::Double", { inputs: { x: 6n } }), + ); + assert.deepEqual( + await wasmEngineModel.executeState("enginedemo::Switch"), + await nativeEngineModel.executeState("enginedemo::Switch"), + ); + const sequenceModel = `package Demo { + private import ScalarValues::*; + metadata def Safety { attribute level : Integer = 2; } + part def Vehicle { attribute mass : Real; } + part seatBelt : Vehicle { @Safety { level = 4; } } + attribute everything [*] = seatBelt.metadata; + }`; + const wasmSequenceModel = await wasm.loads(sequenceModel); + const nativeSequenceModel = await native.loads(sequenceModel); + const wasmSequence = await wasmSequenceModel.eval("Demo::everything"); + assert.equal(wasmSequence.kind, "sequence"); + assert.deepEqual( + wasmSequence, + await nativeSequenceModel.eval("Demo::everything"), + ); + + const multiple = await wasm.parseSources([ + SourceDocument.inline("first.sysml", "package Multi { part def Wheel; }"), + SourceDocument.inline("second.sysml", "package Multi { part car : Wheel; }"), + ]); + assert.ok((await multiple.symbol("Multi::car")).id !== ""); + assert.deepEqual(await multiple.eval("1 + 1"), { kind: "int", value: 2n }); + + await assert.rejects( + () => wasmModel.query(), + (error: unknown) => + error instanceof MissingCapabilityError && error.capability === "query", + ); + await assert.rejects( + () => wasm.rpc.convert({}), + (error: unknown) => error instanceof ConnectError && error.code === Code.Unimplemented, + ); + assert.ok(wasmResponses.includes("GetServerInfo")); + assert.ok(wasmResponses.includes("ParseFile")); + await wasm.close(); + await assert.rejects(() => wasm.loads(SAMPLE), ClosedConnectionError); + }); +} + +test("Node WASM workers terminate on close and stale-version refusal", { skip: wasmSkip }, async () => { + const originalTerminate = Object.getOwnPropertyDescriptor(WorkerPrototype, "terminate") + ?.value as ((this: NodeWorker) => Promise) | undefined; + if (originalTerminate === undefined) { + throw new Error("Worker.prototype.terminate is unavailable"); + } + let termination: Promise | undefined; + WorkerPrototype.terminate = function () { + const pending = originalTerminate.call(this); + termination = pending; + return pending; + }; + try { + const wasmFiles = requireArtifacts(); + const connection = await connectWasm({ + wasm: wasmFiles.wasm, + wasmExec: wasmFiles.wasmExec, + }); + await connection.close(); + assert.ok(connection.isClosed); + await awaitTermination(termination); + await assert.rejects(() => connection.loads(SAMPLE), ClosedConnectionError); + + termination = undefined; + await assert.rejects( + () => + connectWasm({ + wasm: wasmFiles.wasm, + wasmExec: wasmFiles.wasmExec, + version: "nope", + }), + StaleServiceError, + ); + await awaitTermination(termination); + + termination = undefined; + await assert.rejects( + () => + connectWasm({ + wasm: wasmFiles.wasm, + wasmExec: wasmFiles.wasmExec, + requireCapabilities: ["not-a-capability"], + }), + MissingCapabilityError, + ); + await awaitTermination(termination); + } finally { + WorkerPrototype.terminate = originalTerminate; + } +}); + +test("browser inline WASM connections answer client calls", { skip: wasmSkip }, async () => { + const wasmFiles = requireArtifacts(); + const bytes = new Uint8Array(await readFile(wasmFiles.wasm)); + await using connection = await connectBrowserWasm({ + wasm: new Response(bytes, { headers: { "Content-Type": "application/wasm" } }), + wasmExec: pathToFileURL(wasmFiles.wasmExec), + }); + await using native = await connect(); + const wasmModel = await connection.loads(SAMPLE); + const nativeModel = await native.loads(SAMPLE); + assert.equal(wasmModel.hash, nativeModel.hash); + assert.deepEqual(await wasmModel.eval("2 + 2"), await nativeModel.eval("2 + 2")); + assert.equal( + (await wasmModel.symbol("Sample::Car")).id, + (await nativeModel.symbol("Sample::Car")).id, + ); +}); + +const WorkerPrototype = NodeWorker.prototype; + +function requireArtifacts(): WasmArtifacts { + if (artifacts === undefined) { + throw new Error("WASM artifacts are required by this test"); + } + return artifacts; +} + +async function awaitTermination(pending: Promise | undefined): Promise { + if (pending === undefined) { + throw new Error("the WASM worker was not terminated"); + } + await pending; +} + +function clientFor( + host: WasmHost, + options: Parameters[1] = {}, +) { + return createClient(SysMLService, createWasmTransport(host, options)); +} + +function fixedHost(envelope: string): WasmHost { + return { + version: "fake", + call: () => Promise.resolve(envelope), + close: () => Promise.resolve(), + }; +} + +const emptyWasm = new Uint8Array([0, 97, 115, 109, 1, 0, 0, 0]); + +class FakeGo { + readonly importObject = {}; + + run(): Promise { + (globalThis as typeof globalThis & { + sysmlWasm?: { version: string; call(method: string, params: string): string }; + }).sysmlWasm = { + version: "fake-wasm", + call(method, params) { + if (method === "Fail") { + throw new Error("fake failure"); + } + const result = + method === "GetServerInfo" + ? { version: "fake-wasm", capabilities: [] } + : { method, params: JSON.parse(params) as unknown }; + return JSON.stringify({ + jsonrpc: "2.0", + id: null, + result, + }); + }, + }; + return Promise.resolve(); + } +} + +class FakeBrowserWorker { + terminated = false; + + private readonly listeners = new Map(); + private readonly serviceListeners: Array< + (event: MessageEvent) => void + > = []; + + constructor() { + serveWasmPort( + { + postMessage: (message) => { + this.dispatch("message", new MessageEvent("message", { data: message })); + }, + addEventListener: (_type, listener) => this.serviceListeners.push(listener), + start() {}, + }, + { + loadWasm: () => Promise.resolve(emptyWasm), + loadGo: () => Promise.resolve(FakeGo), + }, + ); + } + + addEventListener(type: string, listener: EventListenerOrEventListenerObject): void { + const wrapped = + typeof listener === "function" + ? listener + : (event: Event) => { + listener.handleEvent(event); + }; + const typeListeners = this.listeners.get(type) ?? []; + typeListeners.push(wrapped); + this.listeners.set(type, typeListeners); + } + + postMessage(message: unknown): void { + queueMicrotask(() => { + for (const listener of this.serviceListeners) { + listener({ data: message as WasmWorkerRequest } as MessageEvent); + } + }); + } + + terminate(): Promise { + this.terminated = true; + return Promise.resolve(1); + } + + private dispatch(type: string, event: Event): void { + for (const listener of this.listeners.get(type) ?? []) { + listener(event); + } + } +} + +interface FakePort { + postMessage(message: WasmWorkerRequest): void; + addEventListener( + type: "message", + listener: (event: MessageEvent) => void, + ): void; +} + +function fakePortPair(): { client: FakePort; service: WasmPortLike } { + const clientListeners: Array<(event: MessageEvent) => void> = []; + const serviceListeners: Array<(event: MessageEvent) => void> = []; + const client: FakePort = { + postMessage(message) { + queueMicrotask(() => { + for (const listener of serviceListeners) { + listener({ data: message } as MessageEvent); + } + }); + }, + addEventListener(_type, listener) { + clientListeners.push(listener); + }, + }; + const service: WasmPortLike = { + postMessage(message) { + queueMicrotask(() => { + for (const listener of clientListeners) { + listener({ data: message } as MessageEvent); + } + }); + }, + addEventListener(_type, listener) { + serviceListeners.push(listener); + }, + }; + return { client, service }; +} + +async function waitFor( + messages: WasmWorkerResponse[], + predicate: (message: WasmWorkerResponse) => message is T, +): Promise; +async function waitFor( + messages: WasmWorkerResponse[], + predicate: (message: WasmWorkerResponse) => boolean, +): Promise; +async function waitFor( + messages: WasmWorkerResponse[], + predicate: (message: WasmWorkerResponse) => boolean, +): Promise { + const deadline = Date.now() + 1_000; + while (Date.now() < deadline) { + const match = messages.find(predicate); + if (match !== undefined) { + return match; + } + await new Promise((resolve) => setTimeout(resolve, 0)); + } + throw new Error("timed out waiting for a WASM worker message"); +} + +void (FakeGo satisfies GoConstructor); diff --git a/docs/reference/clients.md b/docs/reference/clients.md index a7b2793051..4af5ac823e 100644 --- a/docs/reference/clients.md +++ b/docs/reference/clients.md @@ -1,15 +1,16 @@ # Client libraries OpenSysML can be reached from a program in seven ways: the Go API, which runs in the calling -process, and six clients of the `sysml-grpc` service. This page describes how to choose between -them, what each covers and what each intentionally leaves out. Each client has an API reference of -its own, and the [client guides](../clients.md) walk through a task with each one. +process, and six clients of the `sysml-grpc` service. The Node client can also run the combined +`sysml-wasm` module directly. This page describes how to choose between them, what each covers +and what each intentionally leaves out. Each client has an API reference of its own, and the +[client guides](../clients.md) walk through a task with each one. | Surface | Reaches the engine by | Published | Full reference | |---|---|---|---| | **Go**, `client/opensysml` | in process; or Connect, to a service someone else runs | with the core (`v*` tags) | [Go packages](api.md) | | **Python**, `opensysml` | gRPC, to a private child service or a named service | PyPI, on the core `v*` tags, at the core's version | [Python API](python-api.md) | -| **Node/TypeScript**, `@openmbee/opensysml` | Connect, to a private child service, a named service, or one a browser page addresses | npm, on core `v*` tags, with per-platform binary packages | [Node API](node-api.md) | +| **Node/TypeScript**, `@openmbee/opensysml` | Connect to a service, or directly to the combined `sysml-wasm` module | npm, on core `v*` tags, with per-platform binary packages | [Node API](node-api.md) | | **Java**, `org.openmbee:opensysml` | Connect, over the JDK's own HTTP client | not on Maven Central; build from a checkout | [Java API](java-api.md) | | **Rust**, `opensysml` | Connect, blocking, no async runtime | crates.io, on core `v*` tags, at the core version | [Rust API](rust-api.md) | | **Julia**, `OpenSysML` | Connect-JSON, over `HTTP.jl` | not in General; develop from a checkout | [Julia API](julia-api.md) | @@ -28,7 +29,8 @@ The protocols and what the service serves on a single port are described in - **In a notebook: Python.** `opensysml` adds generated typed classes, Jupyter display hooks and DataFrame integration to the full RPC surface. - **In a browser or a Node service: `@openmbee/opensysml`.** No native addon, and the browser entry - point needs only `fetch` against a service that allows the page's origin. + point needs only `fetch` against a service that allows the page's origin. Node and browser + callers can also use the combined `sysml-wasm` module without a service. - **In a JVM host application the caller does not control (an Eclipse-based tool, a Cameo plugin, a web application): Java.** Its transport is `java.net.http.HttpClient`, so no gRPC, Netty or `tcnative` dependency reaches the host application. @@ -42,11 +44,13 @@ The protocols and what the service serves on a single port are described in The Go, Java, Julia and MATLAB clients each reach every RPC the service has — the Julia and MATLAB ones through `call`/`callRaw` under the wrapped functions, so nothing on the wire is -out of reach — and so do the Python, Node and Rust clients. +out of reach — and so do the Python, Node's Connect transport and Rust clients. ## What the newer surfaces cover -The Node client covers everything the Python one does, the whole service surface included. +The Node client's Connect transport covers everything the Python one does, the whole service +surface included. Its `connectWasm()` adapter instead exposes the combined module's supported +subset without a service. The Java client covers the whole service surface, as typed immutable results: diff --git a/docs/reference/node-api.md b/docs/reference/node-api.md index 1709b2a920..c97eb7f0f2 100644 --- a/docs/reference/node-api.md +++ b/docs/reference/node-api.md @@ -23,6 +23,8 @@ The package is published on npm. From a checkout, build it with Both re-export the isomorphic core; the browser entry point requires an `address`, since there is nothing to fall back to. +Both also export `connectWasm()` for the combined WebAssembly module; see +[WebAssembly, without a service](../../client/node/README.md#webassembly-without-a-service). ## Opening a connection diff --git a/docs/reference/wasm.md b/docs/reference/wasm.md index ab6b452415..f3f32144fb 100644 --- a/docs/reference/wasm.md +++ b/docs/reference/wasm.md @@ -211,6 +211,9 @@ Measured on a `go1.25` `js/wasm` build: 20,601,256 raw bytes, 5,444,030 bytes wi ## The combined module +The Node and browser client adapter is documented in +[WebAssembly, without a service](../../client/node/README.md#webassembly-without-a-service). + `sysml-wasm` combines the parsing, validation and execution methods of `sysml-core` and `sysml-engine` in one WebAssembly module. It serves `ParseSources`, `ParseFile`, `GetDiagnostics`, `GetSymbol`, `Evaluate`, `Instantiate`, `ExecuteAction`, `ExecuteState` From 42e135dd95e1892afec4761c482312aebec3f9e4 Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Sun, 4 Oct 2026 00:33:04 +0000 Subject: [PATCH 2/5] refactor(node): share the WASM worker host Co-Authored-By: jason.han --- client/node/src/browser/wasm.ts | 159 +++++----------------- client/node/src/core/wasm.ts | 125 ++++++++++++++++- client/node/src/node/wasm-source.ts | 56 ++++++++ client/node/src/node/wasm-worker.ts | 42 +----- client/node/src/node/wasm.ts | 199 ++++------------------------ 5 files changed, 235 insertions(+), 346 deletions(-) create mode 100644 client/node/src/node/wasm-source.ts diff --git a/client/node/src/browser/wasm.ts b/client/node/src/browser/wasm.ts index 6fb7cae9e8..b356cd16d9 100644 --- a/client/node/src/browser/wasm.ts +++ b/client/node/src/browser/wasm.ts @@ -1,12 +1,13 @@ -import { Code, ConnectError } from "@connectrpc/connect"; import { Connection } from "../core/connection.js"; import type { TransportOptions } from "../core/connection.js"; import { connectWasmHost, instantiateInline, loadGoConstructor, + loadWasmSource, + WorkerWasmHost, type WasmHost, - type WasmWorkerRequest, + type WasmWorkerEndpoint, type WasmWorkerResponse, type WorkerWasmSource, } from "../core/wasm.js"; @@ -50,7 +51,7 @@ async function startWorker( wasm: BrowserWasmConnectOptions["wasm"], wasmExec: string | URL | undefined, ): Promise { - const host = new BrowserWasmWorkerHost(worker); + const host = new WorkerWasmHost(browserWorkerEndpoint(worker)); try { const { source, transfer } = await workerSource(wasm); return await host.start( @@ -67,118 +68,28 @@ async function startWorker( } } -class BrowserWasmWorkerHost implements WasmHost { - version = ""; - - private readonly pending = new Map< - number, - { resolve: (envelope: string) => void; reject: (error: unknown) => void } - >(); - private readonly ready: Promise; - private resolveReady!: () => void; - private rejectReady!: (error: unknown) => void; - private nextId = 1; - private failure: ConnectError | undefined; - private closed = false; - - constructor(private readonly worker: Worker) { - this.ready = new Promise((resolve, reject) => { - this.resolveReady = resolve; - this.rejectReady = reject; - }); - worker.addEventListener("message", (event: MessageEvent) => { - this.receive(event.data); - }); - worker.addEventListener("error", (event: ErrorEvent) => { - this.fail(new Error(event.message || "the sysml-wasm worker failed")); - }); - worker.addEventListener("messageerror", () => { - this.fail(new Error("the sysml-wasm worker sent an unreadable message")); - }); - } - - async start( - request: Extract, - transfer: Transferable[], - ): Promise { - this.worker.postMessage(request, transfer); - await this.ready; - return this; - } - - async call(method: string, params: string): Promise { - this.throwIfFailed(); - if (this.closed) { - throw new ConnectError("the sysml-wasm worker is closed", Code.Unavailable); - } - await this.ready; - this.throwIfFailed(); - const id = this.nextId++; - return new Promise((resolve, reject) => { - this.pending.set(id, { resolve, reject }); - try { - this.worker.postMessage({ type: "call", id, method, params }); - } catch (error) { - this.pending.delete(id); - reject(ConnectError.from(error, Code.Unavailable)); - } - }); - } - - close(): Promise { - if (this.closed) { - return Promise.resolve(); - } - this.closed = true; - try { - this.worker.postMessage({ type: "close" }); - } catch { - // A worker that has already failed has no message loop to close. - } - const closed = new ConnectError("the sysml-wasm worker was closed", Code.Unavailable); - for (const pending of this.pending.values()) { - pending.reject(closed); - } - this.pending.clear(); - this.worker.terminate(); - return Promise.resolve(); - } - - private receive(message: WasmWorkerResponse): void { - if (message.type === "ready") { - this.version = message.version; - this.resolveReady(); - return; - } - if (message.type === "failed") { - this.fail(new Error(message.message)); - return; - } - const pending = this.pending.get(message.id); - if (pending !== undefined) { - this.pending.delete(message.id); - pending.resolve(message.envelope); - } - } - - private fail(reason: unknown): void { - if (this.failure !== undefined || this.closed) { - return; - } - this.failure = ConnectError.from(reason, Code.Unavailable); - this.rejectReady(this.failure); - for (const pending of this.pending.values()) { - pending.reject(this.failure); - } - this.pending.clear(); - } - - private throwIfFailed(): void { - const failure = this.failure; - if (failure !== undefined) { - throw failure; - } - } +function browserWorkerEndpoint(worker: Worker): WasmWorkerEndpoint { + return { + post(message, transfer) { + worker.postMessage(message, [...transfer]); + }, + onMessage(listener) { + worker.addEventListener("message", (event: MessageEvent) => { + listener(event.data); + }); + }, + onFailure(listener) { + worker.addEventListener("error", (event: ErrorEvent) => { + listener(new Error(event.message || "the sysml-wasm worker failed")); + }); + worker.addEventListener("messageerror", () => { + listener(new Error("the sysml-wasm worker sent an unreadable message")); + }); + }, + terminate() { + worker.terminate(); + }, + }; } async function workerSource( @@ -216,22 +127,12 @@ async function browserWasmSource( ) { return wasm; } - const response = - wasm instanceof Response ? wasm : await fetch(wasm instanceof URL ? wasm : resourceUrl(wasm)); - if (!response.ok) { - throw new Error(`could not fetch the WebAssembly module: ${response.status}`); - } - if (response.bodyUsed) { + if (wasm instanceof Response && wasm.bodyUsed) { throw new Error("the WebAssembly Response body has already been read"); } - if (typeof WebAssembly.compileStreaming === "function") { - try { - return await WebAssembly.compileStreaming(response.clone()); - } catch { - // Some servers omit the application/wasm content type required by streaming. - } - } - return response.arrayBuffer(); + return loadWasmSource( + wasm instanceof Response ? wasm : wasm instanceof URL ? wasm : resourceUrl(wasm), + ); } function resourceUrl(value: string | URL): string { diff --git a/client/node/src/core/wasm.ts b/client/node/src/core/wasm.ts index 77521f37d0..30bd873770 100644 --- a/client/node/src/core/wasm.ts +++ b/client/node/src/core/wasm.ts @@ -47,6 +47,14 @@ export type WasmWorkerResponse = | { type: "failed"; message: string } | { type: "answer"; id: number; envelope: string }; +/** The platform-specific message and lifecycle surface of a WASM worker. */ +export interface WasmWorkerEndpoint { + post(message: WasmWorkerRequest, transfer: readonly Transferable[]): void; + onMessage(listener: (message: WasmWorkerResponse) => void): void; + onFailure(listener: (reason: Error) => void): void; + terminate(): Promise | void; +} + /** The message port surface shared by browser and Node workers. */ export interface WasmPortLike { postMessage(message: WasmWorkerResponse): void; @@ -256,6 +264,117 @@ export async function connectWasmHost( }); } +/** Adapts a platform worker endpoint to the shared WASM host protocol. */ +export class WorkerWasmHost implements WasmHost { + version = ""; + + private readonly pending = new Map< + number, + { resolve: (envelope: string) => void; reject: (error: unknown) => void } + >(); + private readonly ready: Promise; + private resolveReady!: () => void; + private rejectReady!: (error: unknown) => void; + private nextId = 1; + private failure: ConnectError | undefined; + private closed = false; + + constructor(private readonly endpoint: WasmWorkerEndpoint) { + this.ready = new Promise((resolve, reject) => { + this.resolveReady = resolve; + this.rejectReady = reject; + }); + endpoint.onMessage((message) => { + this.receive(message); + }); + endpoint.onFailure((reason) => { + this.fail(reason); + }); + } + + async start( + request: Extract, + transfer: readonly Transferable[], + ): Promise { + this.endpoint.post(request, transfer); + await this.ready; + return this; + } + + async call(method: string, params: string): Promise { + this.throwIfFailed(); + if (this.closed) { + throw new ConnectError("the sysml-wasm worker is closed", Code.Unavailable); + } + await this.ready; + this.throwIfFailed(); + const id = this.nextId++; + return new Promise((resolve, reject) => { + this.pending.set(id, { resolve, reject }); + try { + this.endpoint.post({ type: "call", id, method, params }, []); + } catch (error) { + this.pending.delete(id); + reject(ConnectError.from(error, Code.Unavailable)); + } + }); + } + + async close(): Promise { + if (this.closed) { + return; + } + this.closed = true; + try { + this.endpoint.post({ type: "close" }, []); + } catch { + // A worker that has already failed has no message loop to close. + } + const closed = new ConnectError("the sysml-wasm worker was closed", Code.Unavailable); + for (const pending of this.pending.values()) { + pending.reject(closed); + } + this.pending.clear(); + await this.endpoint.terminate(); + } + + private receive(message: WasmWorkerResponse): void { + if (message.type === "ready") { + this.version = message.version; + this.resolveReady(); + return; + } + if (message.type === "failed") { + this.fail(new Error(message.message)); + return; + } + const pending = this.pending.get(message.id); + if (pending !== undefined) { + this.pending.delete(message.id); + pending.resolve(message.envelope); + } + } + + private fail(reason: Error): void { + if (this.failure !== undefined || this.closed) { + return; + } + this.failure = ConnectError.from(reason, Code.Unavailable); + this.rejectReady(this.failure); + for (const pending of this.pending.values()) { + pending.reject(this.failure); + } + this.pending.clear(); + } + + private throwIfFailed(): void { + const failure = this.failure; + if (failure !== undefined) { + throw failure; + } + } +} + /** * Serves the shared init/call protocol on a Node or browser worker port. */ @@ -431,8 +550,8 @@ function errorEnvelope(error: unknown): string { }); } -async function loadWasmSource( - source: WorkerWasmSource, +export async function loadWasmSource( + source: WorkerWasmSource | URL | Response, ): Promise { if ( source instanceof WebAssembly.Module || @@ -441,7 +560,7 @@ async function loadWasmSource( ) { return source; } - const response = await fetch(source); + const response = source instanceof Response ? source : await fetch(source); if (!response.ok) { throw new OpenSysMLError(`could not fetch the WebAssembly module: ${response.status}`); } diff --git a/client/node/src/node/wasm-source.ts b/client/node/src/node/wasm-source.ts new file mode 100644 index 0000000000..71941b2552 --- /dev/null +++ b/client/node/src/node/wasm-source.ts @@ -0,0 +1,56 @@ +import { readFile } from "node:fs/promises"; +import { isAbsolute, resolve } from "node:path"; +import { pathToFileURL } from "node:url"; +import type { WorkerWasmSource } from "../core/wasm.js"; + +export async function loadNodeWasm( + source: WorkerWasmSource | URL, +): Promise { + if ( + source instanceof WebAssembly.Module || + source instanceof ArrayBuffer || + ArrayBuffer.isView(source) + ) { + return source; + } + const url = source instanceof URL ? source : parseUrl(source); + if (url?.protocol === "file:") { + return readFile(url); + } + if (url !== undefined) { + const response = await fetch(url); + if (!response.ok) { + throw new Error(`could not fetch the WebAssembly module: ${response.status}`); + } + return response.arrayBuffer(); + } + if (typeof source === "string") { + return readFile(isAbsolute(source) ? source : resolve(source)); + } + return readFile(source); +} + +export function moduleSpecifier(source: string | URL): string { + if (source instanceof URL) { + return source.href; + } + if (!/^[A-Za-z]:[\\/]/.test(source)) { + try { + return new URL(source).href; + } catch { + // Resolve filesystem paths relative to the current working directory. + } + } + return pathToFileURL(isAbsolute(source) ? source : resolve(source)).href; +} + +export function parseUrl(source: string): URL | undefined { + if (/^[A-Za-z]:[\\/]/.test(source)) { + return undefined; + } + try { + return new URL(source); + } catch { + return undefined; + } +} diff --git a/client/node/src/node/wasm-worker.ts b/client/node/src/node/wasm-worker.ts index 4a68ed8761..7d8beaecd7 100644 --- a/client/node/src/node/wasm-worker.ts +++ b/client/node/src/node/wasm-worker.ts @@ -1,13 +1,10 @@ -import { readFile } from "node:fs/promises"; -import { resolve } from "node:path"; -import { pathToFileURL } from "node:url"; import { parentPort } from "node:worker_threads"; import { loadGoConstructor, serveWasmPort, type WasmPortLike, - type WorkerWasmSource, } from "../core/wasm.js"; +import { loadNodeWasm, moduleSpecifier } from "./wasm-source.js"; if (parentPort === null) { throw new Error("the sysml-wasm worker must run in a worker thread"); @@ -17,40 +14,3 @@ serveWasmPort(parentPort as unknown as WasmPortLike, { loadWasm: loadNodeWasm, loadGo: (wasmExec) => loadGoConstructor(wasmExec === undefined ? undefined : moduleSpecifier(wasmExec)), }); - -async function loadNodeWasm(source: WorkerWasmSource): Promise { - if ( - source instanceof WebAssembly.Module || - source instanceof ArrayBuffer || - ArrayBuffer.isView(source) - ) { - return source; - } - const url = parseUrl(source); - if (url?.protocol === "file:") { - return readFile(url); - } - if (url !== undefined) { - const response = await fetch(url); - if (!response.ok) { - throw new Error(`could not fetch the WebAssembly module: ${response.status}`); - } - return response.arrayBuffer(); - } - return readFile(resolve(source)); -} - -function moduleSpecifier(source: string): string { - return parseUrl(source)?.href ?? pathToFileURL(resolve(source)).href; -} - -function parseUrl(value: string): URL | undefined { - if (/^[A-Za-z]:[\\/]/.test(value)) { - return undefined; - } - try { - return new URL(value); - } catch { - return undefined; - } -} diff --git a/client/node/src/node/wasm.ts b/client/node/src/node/wasm.ts index 9ce87541ae..c386e01f6e 100644 --- a/client/node/src/node/wasm.ts +++ b/client/node/src/node/wasm.ts @@ -1,19 +1,16 @@ -import { readFile } from "node:fs/promises"; -import { isAbsolute, resolve } from "node:path"; -import { pathToFileURL } from "node:url"; import { Worker } from "node:worker_threads"; -import { Code, ConnectError } from "@connectrpc/connect"; import { Connection } from "../core/connection.js"; import type { TransportOptions } from "../core/connection.js"; import { connectWasmHost, instantiateInline, loadGoConstructor, + WorkerWasmHost, type WasmHost, - type WasmWorkerRequest, - type WasmWorkerResponse, + type WasmWorkerEndpoint, type WorkerWasmSource, } from "../core/wasm.js"; +import { loadNodeWasm, moduleSpecifier } from "./wasm-source.js"; /** Options for connecting to a Go WebAssembly module from Node. */ export interface WasmConnectOptions extends TransportOptions { @@ -45,7 +42,7 @@ export async function connectWasm(options: WasmConnectOptions): Promise { const worker = new Worker(new URL("./wasm-worker.js", import.meta.url)); - const host = new NodeWasmWorkerHost(worker); + const host = new WorkerWasmHost(nodeWorkerEndpoint(worker)); const { source, transfer } = workerSource(wasm); try { return await host.start( @@ -58,121 +55,6 @@ async function startWorker(wasm: WasmConnectOptions["wasm"], wasmExec: string): } } -class NodeWasmWorkerHost implements WasmHost { - version = ""; - - private readonly pending = new Map< - number, - { resolve: (envelope: string) => void; reject: (error: unknown) => void } - >(); - private readonly ready: Promise; - private resolveReady!: () => void; - private rejectReady!: (error: unknown) => void; - private nextId = 1; - private failure: ConnectError | undefined; - private closed = false; - - constructor(private readonly worker: Worker) { - this.ready = new Promise((resolve, reject) => { - this.resolveReady = resolve; - this.rejectReady = reject; - }); - worker.on("message", (message: WasmWorkerResponse) => { - this.receive(message); - }); - worker.on("error", (error: Error) => { - this.fail(error); - }); - worker.on("exit", (code: number) => { - if (!this.closed) { - this.fail(new Error(`the sysml-wasm worker exited with code ${code}`)); - } - }); - } - - async start( - request: Extract, - transfer: readonly ArrayBuffer[], - ): Promise { - this.worker.postMessage(request, transfer); - await this.ready; - return this; - } - - async call(method: string, params: string): Promise { - this.throwIfFailed(); - if (this.closed) { - throw new ConnectError("the sysml-wasm worker is closed", Code.Unavailable); - } - await this.ready; - this.throwIfFailed(); - const id = this.nextId++; - return new Promise((resolve, reject) => { - this.pending.set(id, { resolve, reject }); - try { - this.worker.postMessage({ type: "call", id, method, params }); - } catch (error) { - this.pending.delete(id); - reject(ConnectError.from(error, Code.Unavailable)); - } - }); - } - - async close(): Promise { - if (this.closed) { - return; - } - this.closed = true; - try { - this.worker.postMessage({ type: "close" }); - } catch { - // A worker that has already failed has no message loop to close. - } - const closed = new ConnectError("the sysml-wasm worker was closed", Code.Unavailable); - for (const pending of this.pending.values()) { - pending.reject(closed); - } - this.pending.clear(); - await this.worker.terminate(); - } - - private receive(message: WasmWorkerResponse): void { - if (message.type === "ready") { - this.version = message.version; - this.resolveReady(); - return; - } - if (message.type === "failed") { - this.fail(new Error(message.message)); - return; - } - const pending = this.pending.get(message.id); - if (pending !== undefined) { - this.pending.delete(message.id); - pending.resolve(message.envelope); - } - } - - private fail(reason: unknown): void { - if (this.failure !== undefined || this.closed) { - return; - } - this.failure = ConnectError.from(reason, Code.Unavailable); - this.rejectReady(this.failure); - for (const pending of this.pending.values()) { - pending.reject(this.failure); - } - this.pending.clear(); - } - - private throwIfFailed(): void { - const failure = this.failure; - if (failure !== undefined) { - throw failure; - } - } -} - function workerSource( wasm: WasmConnectOptions["wasm"], ): { source: WorkerWasmSource; transfer: readonly ArrayBuffer[] } { @@ -183,55 +65,26 @@ function workerSource( return { source: wasm instanceof URL ? wasm.href : wasm, transfer: [] }; } -async function loadNodeWasm(wasm: WasmConnectOptions["wasm"]): Promise { - if (wasm instanceof Uint8Array) { - return wasm; - } - if (wasm instanceof URL) { - if (wasm.protocol === "file:") { - return readFile(wasm); - } - const response = await fetch(wasm); - if (!response.ok) { - throw new Error(`could not fetch the WebAssembly module: ${response.status}`); - } - return response.arrayBuffer(); - } - const url = parseUrl(wasm); - if (url?.protocol === "file:") { - return readFile(url); - } - if (url !== undefined) { - const response = await fetch(url); - if (!response.ok) { - throw new Error(`could not fetch the WebAssembly module: ${response.status}`); - } - return response.arrayBuffer(); - } - return readFile(wasm); -} - -function moduleSpecifier(source: string | URL): string { - if (source instanceof URL) { - return source.href; - } - if (!/^[A-Za-z]:[\\/]/.test(source)) { - try { - return new URL(source).href; - } catch { - // Resolve filesystem paths relative to the current working directory. - } - } - return pathToFileURL(isAbsolute(source) ? source : resolve(source)).href; -} - -function parseUrl(source: string): URL | undefined { - if (/^[A-Za-z]:[\\/]/.test(source)) { - return undefined; - } - try { - return new URL(source); - } catch { - return undefined; - } +function nodeWorkerEndpoint(worker: Worker): WasmWorkerEndpoint { + let closed = false; + return { + post(message, transfer) { + worker.postMessage(message, [...transfer] as ArrayBuffer[]); + }, + onMessage(listener) { + worker.on("message", listener); + }, + onFailure(listener) { + worker.on("error", listener); + worker.on("exit", (code: number) => { + if (!closed) { + listener(new Error(`the sysml-wasm worker exited with code ${code}`)); + } + }); + }, + async terminate() { + closed = true; + await worker.terminate(); + }, + }; } From b9eddc851e6b27a3d93a76a65fa436bc91d7d807 Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Sun, 4 Oct 2026 01:21:51 +0000 Subject: [PATCH 3/5] fix(node): cancel abandoned WASM worker calls and load each Go runtime once Co-Authored-By: jason.han --- client/node/src/core/wasm.ts | 140 +++++++++++++++++++++++++++------- client/node/src/node/wasm.ts | 2 +- client/node/test/wasm.test.ts | 135 ++++++++++++++++++++++++++++++++ 3 files changed, 247 insertions(+), 30 deletions(-) diff --git a/client/node/src/core/wasm.ts b/client/node/src/core/wasm.ts index 30bd873770..5757a0a06a 100644 --- a/client/node/src/core/wasm.ts +++ b/client/node/src/core/wasm.ts @@ -22,7 +22,7 @@ import { encodingOf, interceptors, timeoutOf } from "./transport.js"; /** A host surface implemented by an inline Go runtime or a worker. */ export interface WasmHost { readonly version: string; - call(method: string, params: string): Promise; + call(method: string, params: string, signal?: AbortSignal): Promise; close(): Promise; } @@ -35,6 +35,9 @@ export interface GoRuntime { /** A constructor exported globally by wasm_exec.js. */ export type GoConstructor = new () => GoRuntime; +const goConstructors = new Map>(); +let goConstructorLoads = Promise.resolve(); + export type WorkerWasmSource = string | ArrayBuffer | Uint8Array | WebAssembly.Module; export type WasmWorkerRequest = @@ -192,7 +195,7 @@ export function createWasmTransport(host: WasmHost, options: TransportOptions): throw abortError(request.signal); } const params = JSON.stringify(toJson(method.input, request.message)); - const answer = host.call(method.name, params); + const answer = host.call(method.name, params, request.signal); const envelope = await withAbort(answer, request.signal); const body = readEnvelope(envelope); return { @@ -268,10 +271,7 @@ export async function connectWasmHost( export class WorkerWasmHost implements WasmHost { version = ""; - private readonly pending = new Map< - number, - { resolve: (envelope: string) => void; reject: (error: unknown) => void } - >(); + private readonly pending = new Map(); private readonly ready: Promise; private resolveReady!: () => void; private rejectReady!: (error: unknown) => void; @@ -301,21 +301,43 @@ export class WorkerWasmHost implements WasmHost { return this; } - async call(method: string, params: string): Promise { + async call(method: string, params: string, signal?: AbortSignal): Promise { + if (signal?.aborted) { + throw abortError(signal); + } this.throwIfFailed(); if (this.closed) { throw new ConnectError("the sysml-wasm worker is closed", Code.Unavailable); } - await this.ready; + if (signal === undefined) { + await this.ready; + } else { + await withAbort(this.ready, signal); + } + if (signal?.aborted) { + throw abortError(signal); + } this.throwIfFailed(); + this.throwIfClosed(); const id = this.nextId++; return new Promise((resolve, reject) => { - this.pending.set(id, { resolve, reject }); + if (signal === undefined) { + this.pending.set(id, { resolve, reject }); + } else { + const onAbort = (): void => { + this.removePending(id)?.reject(abortError(signal)); + }; + this.pending.set(id, { resolve, reject, signal, onAbort }); + signal.addEventListener("abort", onAbort, { once: true }); + if (signal.aborted) { + onAbort(); + return; + } + } try { this.endpoint.post({ type: "call", id, method, params }, []); } catch (error) { - this.pending.delete(id); - reject(ConnectError.from(error, Code.Unavailable)); + this.removePending(id)?.reject(ConnectError.from(error, Code.Unavailable)); } }); } @@ -331,10 +353,9 @@ export class WorkerWasmHost implements WasmHost { // A worker that has already failed has no message loop to close. } const closed = new ConnectError("the sysml-wasm worker was closed", Code.Unavailable); - for (const pending of this.pending.values()) { - pending.reject(closed); + for (const id of this.pending.keys()) { + this.removePending(id)?.reject(closed); } - this.pending.clear(); await this.endpoint.terminate(); } @@ -348,9 +369,8 @@ export class WorkerWasmHost implements WasmHost { this.fail(new Error(message.message)); return; } - const pending = this.pending.get(message.id); + const pending = this.removePending(message.id); if (pending !== undefined) { - this.pending.delete(message.id); pending.resolve(message.envelope); } } @@ -361,10 +381,20 @@ export class WorkerWasmHost implements WasmHost { } this.failure = ConnectError.from(reason, Code.Unavailable); this.rejectReady(this.failure); - for (const pending of this.pending.values()) { - pending.reject(this.failure); + for (const id of this.pending.keys()) { + this.removePending(id)?.reject(this.failure); + } + } + + private removePending(id: number): PendingWorkerCall | undefined { + const pending = this.pending.get(id); + if (pending !== undefined) { + this.pending.delete(id); + if (pending.signal !== undefined && pending.onAbort !== undefined) { + pending.signal.removeEventListener("abort", pending.onAbort); + } } - this.pending.clear(); + return pending; } private throwIfFailed(): void { @@ -373,6 +403,19 @@ export class WorkerWasmHost implements WasmHost { throw failure; } } + + private throwIfClosed(): void { + if (this.closed) { + throw new ConnectError("the sysml-wasm worker is closed", Code.Unavailable); + } + } +} + +interface PendingWorkerCall { + resolve: (envelope: string) => void; + reject: (error: unknown) => void; + signal?: AbortSignal; + onAbort?: () => void; } /** @@ -443,13 +486,41 @@ export function serveWasmPort(port: WasmPortLike, loaders: WasmPortLoaders = {}) port.start?.(); } -/** Loads the Go constructor installed by wasm_exec.js, importing it when needed. */ -export async function loadGoConstructor( - wasmExec?: string, - forceImport = false, -): Promise { +/** Loads and caches the Go constructor installed by wasm_exec.js. */ +export function loadGoConstructor(wasmExec?: string): Promise { const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; - if (wasmExec !== undefined && (forceImport || runtime.Go === undefined)) { + if (wasmExec === undefined) { + if (runtime.Go === undefined) { + return Promise.reject( + new OpenSysMLError("Go is unavailable; supply the matching wasm_exec.js"), + ); + } + return Promise.resolve(runtime.Go); + } + const cached = goConstructors.get(wasmExec); + if (cached !== undefined) { + return cached; + } + const loading = goConstructorLoads.then(() => loadGoConstructorFromScript(wasmExec)); + goConstructors.set(wasmExec, loading); + goConstructorLoads = loading.then( + () => undefined, + () => undefined, + ); + void loading.catch(() => { + if (goConstructors.get(wasmExec) === loading) { + goConstructors.delete(wasmExec); + } + }); + return loading; +} + +async function loadGoConstructorFromScript(wasmExec: string): Promise { + const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; + const previous = runtime.Go; + delete runtime.Go; + let loaded: GoConstructor | undefined; + try { try { await import(wasmExec); } catch (cause) { @@ -461,11 +532,22 @@ export async function loadGoConstructor( } importScripts(wasmExec); } + loaded = globalGoConstructor(); + if (loaded === undefined) { + throw new OpenSysMLError("Go is unavailable; supply the matching wasm_exec.js"); + } + return loaded; + } finally { + if (previous !== undefined) { + runtime.Go = previous; + } else if (loaded === undefined) { + delete runtime.Go; + } } - if (runtime.Go === undefined) { - throw new OpenSysMLError("Go is unavailable; supply the matching wasm_exec.js"); - } - return runtime.Go; +} + +function globalGoConstructor(): GoConstructor | undefined { + return (globalThis as typeof globalThis & { Go?: GoConstructor }).Go; } function readEnvelope(envelope: string): unknown { diff --git a/client/node/src/node/wasm.ts b/client/node/src/node/wasm.ts index c386e01f6e..65cd87cb58 100644 --- a/client/node/src/node/wasm.ts +++ b/client/node/src/node/wasm.ts @@ -32,7 +32,7 @@ export async function connectWasm(options: WasmConnectOptions): Promise { + const directory = mkdtempSync(join(tmpdir(), "opensysml-wasm-exec-")); + const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; + const initial = runtime.Go; + const pathA = join(directory, "go-a.mjs"); + const pathB = join(directory, "go-b.mjs"); + writeFileSync(pathA, "globalThis.Go = class GoA {};\n"); + writeFileSync(pathB, "globalThis.Go = class GoB {};\n"); + runtime.Go = FakeGo; + + try { + const goA = await loadGoConstructor(pathToFileURL(pathA).href); + const goB = await loadGoConstructor(pathToFileURL(pathB).href); + const goAAgain = await loadGoConstructor(pathToFileURL(pathA).href); + + assert.equal(goA.name, "GoA"); + assert.equal(goB.name, "GoB"); + assert.strictEqual(goAAgain, goA); + assert.strictEqual(runtime.Go, FakeGo); + } finally { + if (initial === undefined) { + delete runtime.Go; + } else { + runtime.Go = initial; + } + rmSync(directory, { recursive: true, force: true }); + } +}); + +test("Go constructor load failures are not cached", async () => { + const directory = mkdtempSync(join(tmpdir(), "opensysml-wasm-exec-")); + const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; + const initial = runtime.Go; + const path = join(directory, "go-retry.mjs"); + const specifier = `${pathToFileURL(path).href}?missing`; + writeFileSync(path, "export {};\n"); + delete runtime.Go; + + try { + const firstFailure = loadGoConstructor(specifier); + await assert.rejects(firstFailure, /Go is unavailable/); + const retry = loadGoConstructor(specifier); + assert.notStrictEqual(retry, firstFailure); + await assert.rejects(retry, /Go is unavailable/); + + writeFileSync(path, "globalThis.Go = class GoRetry {};\n"); + const GoRetry = await loadGoConstructor(`${pathToFileURL(path).href}?retry`); + assert.equal(GoRetry.name, "GoRetry"); + } finally { + if (initial === undefined) { + delete runtime.Go; + } else { + runtime.Go = initial; + } + rmSync(directory, { recursive: true, force: true }); + } +}); + +test("worker calls remove aborted pending entries and ignore late answers", async () => { + const posted: WasmWorkerRequest[] = []; + let onMessage: ((message: WasmWorkerResponse) => void) | undefined; + const endpoint: WasmWorkerEndpoint = { + post(message) { + posted.push(message); + if (message.type === "init") { + onMessage?.({ type: "ready", version: "fake-worker" }); + } + }, + onMessage(listener) { + onMessage = listener; + }, + onFailure() {}, + terminate() {}, + }; + const host = new WorkerWasmHost(endpoint); + await host.start({ type: "init", wasm: emptyWasm }, []); + + const controllers = Array.from({ length: 100 }, () => new AbortController()); + const calls = controllers.map((controller) => + host.call("NeverAnswers", "{}", controller.signal), + ); + await new Promise((resolve) => setImmediate(resolve)); + const requests = posted.filter( + (message): message is Extract => + message.type === "call", + ); + assert.equal(requests.length, 100); + + for (const controller of controllers) { + controller.abort(new ConnectError("cancelled", Code.Canceled)); + } + const results = await Promise.allSettled(calls); + assert.ok( + results.every( + (result) => + result.status === "rejected" && + result.reason instanceof ConnectError && + result.reason.code === Code.Canceled, + ), + ); + const pending = (host as unknown as { pending: Map }).pending; + assert.equal(pending.size, 0); + + const alreadyAborted = new AbortController(); + alreadyAborted.abort(new ConnectError("cancelled", Code.Canceled)); + const postedBeforeAbortedCall = posted.length; + await assert.rejects(host.call("AlreadyAborted", "{}", alreadyAborted.signal), { + code: Code.Canceled, + }); + assert.equal(posted.length, postedBeforeAbortedCall); + + const firstRequest = requests[0]; + assert.ok(firstRequest); + onMessage?.({ type: "answer", id: firstRequest.id, envelope: "late answer" }); + assert.equal(pending.size, 0); + + const normalCall = host.call("AfterCancel", "{}"); + await new Promise((resolve) => setImmediate(resolve)); + const normalRequest = posted.find( + (message): message is Extract => + message.type === "call" && message.method === "AfterCancel", + ); + assert.ok(normalRequest); + onMessage?.({ type: "answer", id: normalRequest.id, envelope: "normal answer" }); + assert.equal(await normalCall, "normal answer"); + assert.equal(pending.size, 0); + + await host.close(); +}); + test("browser WASM workers connect through the shared request protocol", async () => { const sources = [ emptyWasm, From 3f90f62cffda83f06ba3be2178e2e213ac9c4e36 Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Sun, 4 Oct 2026 01:44:34 +0000 Subject: [PATCH 4/5] fix(node): keep the shared Go global while loading a runtime Co-Authored-By: jason.han --- client/node/src/core/wasm.ts | 24 +++++++-------- client/node/test/wasm.test.ts | 55 +++++++++++++++++++++++++++++++++++ 2 files changed, 66 insertions(+), 13 deletions(-) diff --git a/client/node/src/core/wasm.ts b/client/node/src/core/wasm.ts index 5757a0a06a..d340c1dfa0 100644 --- a/client/node/src/core/wasm.ts +++ b/client/node/src/core/wasm.ts @@ -488,14 +488,14 @@ export function serveWasmPort(port: WasmPortLike, loaders: WasmPortLoaders = {}) /** Loads and caches the Go constructor installed by wasm_exec.js. */ export function loadGoConstructor(wasmExec?: string): Promise { - const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; if (wasmExec === undefined) { - if (runtime.Go === undefined) { - return Promise.reject( - new OpenSysMLError("Go is unavailable; supply the matching wasm_exec.js"), - ); - } - return Promise.resolve(runtime.Go); + return goConstructorLoads.then(() => { + const Go = globalGoConstructor(); + if (Go === undefined) { + throw new OpenSysMLError("Go is unavailable; supply the matching wasm_exec.js"); + } + return Go; + }); } const cached = goConstructors.get(wasmExec); if (cached !== undefined) { @@ -518,8 +518,6 @@ export function loadGoConstructor(wasmExec?: string): Promise { async function loadGoConstructorFromScript(wasmExec: string): Promise { const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; const previous = runtime.Go; - delete runtime.Go; - let loaded: GoConstructor | undefined; try { try { await import(wasmExec); @@ -532,16 +530,16 @@ async function loadGoConstructorFromScript(wasmExec: string): Promise { + const directory = mkdtempSync(join(tmpdir(), "opensysml-wasm-exec-")); + const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; + const initial = runtime.Go; + const path = join(directory, "go-preloaded.mjs"); + const specifier = `${pathToFileURL(path).href}?preloaded`; + writeFileSync(path, "globalThis.Go = class GoPreloaded {};\n"); + + try { + await import(specifier); + const preloaded = runtime.Go; + assert.ok(preloaded); + assert.equal(preloaded.name, "GoPreloaded"); + assert.strictEqual(await loadGoConstructor(specifier), preloaded); + } finally { + if (initial === undefined) { + delete runtime.Go; + } else { + runtime.Go = initial; + } + rmSync(directory, { recursive: true, force: true }); + } +}); + +test("Go constructor lookups without a runtime wait for explicit loads", async () => { + const directory = mkdtempSync(join(tmpdir(), "opensysml-wasm-exec-")); + const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; + const initial = runtime.Go; + const path = join(directory, "go-slow.mjs"); + const specifier = pathToFileURL(path).href; + const preloaded = class PreloadedGo extends FakeGo {}; + runtime.Go = preloaded; + writeFileSync( + path, + "await new Promise((resolve) => setTimeout(resolve, 50));\n" + + "globalThis.Go = class GoSlow {};\n", + ); + const slowLoad = loadGoConstructor(specifier); + const queuedLookup = loadGoConstructor(); + + try { + assert.strictEqual(await queuedLookup, preloaded); + assert.equal((await slowLoad).name, "GoSlow"); + assert.strictEqual(runtime.Go, preloaded); + } finally { + await Promise.allSettled([slowLoad, queuedLookup]); + if (initial === undefined) { + delete runtime.Go; + } else { + runtime.Go = initial; + } + rmSync(directory, { recursive: true, force: true }); + } +}); + test("Go constructor load failures are not cached", async () => { const directory = mkdtempSync(join(tmpdir(), "opensysml-wasm-exec-")); const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; From 592e0faa5a4ffb3a2c7767f8368119402061c29c Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Sun, 4 Oct 2026 01:49:40 +0000 Subject: [PATCH 5/5] fix(node): accept only a Go runtime its script installed Co-Authored-By: jason.han --- client/node/README.md | 3 ++ client/node/src/core/wasm.ts | 25 ++++++++--- client/node/test/wasm.test.ts | 80 +++++++++++++++++++++++++++++++++-- 3 files changed, 99 insertions(+), 9 deletions(-) diff --git a/client/node/README.md b/client/node/README.md index 4d3dce99c4..8d80c57694 100644 --- a/client/node/README.md +++ b/client/node/README.md @@ -218,6 +218,9 @@ await using connection = await connectWasm({ }); ``` +If the page loads `wasm_exec.js` itself, omit `wasmExec` to use the installed Go +constructor. + The browser worker module can be bundled from `@openmbee/opensysml/browser/wasm-worker`. The combined module measures about 7.8 MB gzipped and 5.5 MB with Brotli. A package containing the matching WASM diff --git a/client/node/src/core/wasm.ts b/client/node/src/core/wasm.ts index d340c1dfa0..e1a7828f5f 100644 --- a/client/node/src/core/wasm.ts +++ b/client/node/src/core/wasm.ts @@ -518,6 +518,7 @@ export function loadGoConstructor(wasmExec?: string): Promise { async function loadGoConstructorFromScript(wasmExec: string): Promise { const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; const previous = runtime.Go; + let loaded: GoConstructor | undefined; try { try { await import(wasmExec); @@ -530,16 +531,28 @@ async function loadGoConstructorFromScript(wasmExec: string): Promise { +test("preloaded Go runtime scripts require an omitted wasmExec", async () => { const directory = mkdtempSync(join(tmpdir(), "opensysml-wasm-exec-")); const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; const initial = runtime.Go; @@ -218,7 +218,49 @@ test("preloaded Go runtime scripts are captured from the global", async () => { const preloaded = runtime.Go; assert.ok(preloaded); assert.equal(preloaded.name, "GoPreloaded"); - assert.strictEqual(await loadGoConstructor(specifier), preloaded); + const message = `wasm_exec.js at ${specifier} did not install Go here; it may already have run in this realm. Omit wasmExec to use the installed Go, or load the runtime in a worker`; + await assert.rejects( + loadGoConstructor(specifier), + (error: unknown) => error instanceof OpenSysMLError && error.message === message, + ); + assert.strictEqual(runtime.Go, preloaded); + assert.strictEqual(await loadGoConstructor(), preloaded); + } finally { + if (initial === undefined) { + delete runtime.Go; + } else { + runtime.Go = initial; + } + rmSync(directory, { recursive: true, force: true }); + } +}); + +test("empty Go runtime scripts are rejected without caching an existing constructor", async () => { + const directory = mkdtempSync(join(tmpdir(), "opensysml-wasm-exec-")); + const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; + const initial = runtime.Go; + const path = join(directory, "go-empty.mjs"); + const specifier = `${pathToFileURL(path).href}?empty`; + const previous = class PreviousGo extends FakeGo {}; + const message = `wasm_exec.js at ${specifier} did not install Go here; it may already have run in this realm. Omit wasmExec to use the installed Go, or load the runtime in a worker`; + writeFileSync(path, "export {};\n"); + runtime.Go = previous; + + try { + const firstFailure = loadGoConstructor(specifier); + await assert.rejects( + firstFailure, + (error: unknown) => error instanceof OpenSysMLError && error.message === message, + ); + assert.strictEqual(runtime.Go, previous); + + const retry = loadGoConstructor(specifier); + assert.notStrictEqual(retry, firstFailure); + await assert.rejects( + retry, + (error: unknown) => error instanceof OpenSysMLError && error.message === message, + ); + assert.strictEqual(runtime.Go, previous); } finally { if (initial === undefined) { delete runtime.Go; @@ -260,6 +302,28 @@ test("Go constructor lookups without a runtime wait for explicit loads", async ( } }); +test("failed Go runtime scripts restore an initially empty global", async () => { + const directory = mkdtempSync(join(tmpdir(), "opensysml-wasm-exec-")); + const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; + const initial = runtime.Go; + const path = join(directory, "go-broken.mjs"); + const specifier = `${pathToFileURL(path).href}?broken`; + writeFileSync(path, 'globalThis.Go = class GoBroken {};\nthrow new Error("failed");\n'); + delete runtime.Go; + + try { + await assert.rejects(loadGoConstructor(specifier), /could not load wasm_exec\.js/); + assert.equal(runtime.Go, undefined); + } finally { + if (initial === undefined) { + delete runtime.Go; + } else { + runtime.Go = initial; + } + rmSync(directory, { recursive: true, force: true }); + } +}); + test("Go constructor load failures are not cached", async () => { const directory = mkdtempSync(join(tmpdir(), "opensysml-wasm-exec-")); const runtime = globalThis as typeof globalThis & { Go?: GoConstructor }; @@ -429,6 +493,16 @@ for (const thread of ["worker", "inline"] as const) { const nativeModel = await native.loads(SAMPLE); assert.equal(wasmModel.hash, nativeModel.hash); assert.deepEqual(wasmModel.diagnostics, nativeModel.diagnostics); + if (thread === "inline") { + await using secondWasm = await connectWasm({ + wasm: wasmFiles.wasm, + wasmExec: wasmFiles.wasmExec, + thread: "inline", + }); + const secondModel = await secondWasm.loads(SAMPLE); + assert.equal(secondModel.hash, nativeModel.hash); + assert.deepEqual(await secondModel.eval("2 + 2"), await nativeModel.eval("2 + 2")); + } const wasmCar = await wasmModel.symbol("Sample::Car"); const nativeCar = await nativeModel.symbol("Sample::Car"); assert.equal(wasmCar.id, nativeCar.id); @@ -580,9 +654,9 @@ test("Node WASM workers terminate on close and stale-version refusal", { skip: w test("browser inline WASM connections answer client calls", { skip: wasmSkip }, async () => { const wasmFiles = requireArtifacts(); const bytes = new Uint8Array(await readFile(wasmFiles.wasm)); + await import(pathToFileURL(wasmFiles.wasmExec).href); await using connection = await connectBrowserWasm({ wasm: new Response(bytes, { headers: { "Content-Type": "application/wasm" } }), - wasmExec: pathToFileURL(wasmFiles.wasmExec), }); await using native = await connect(); const wasmModel = await connection.loads(SAMPLE);