diff --git a/apps/docs/src/content/docs/guides/handle-phases-and-errors.mdx b/apps/docs/src/content/docs/guides/handle-phases-and-errors.mdx index 47a4e7c6..33fd0639 100644 --- a/apps/docs/src/content/docs/guides/handle-phases-and-errors.mdx +++ b/apps/docs/src/content/docs/guides/handle-phases-and-errors.mdx @@ -126,7 +126,9 @@ Every error the Wavelength SDK originates extends [`WavelengthError`](/reference string that identifies the failure reason at the machine level. Match on `err.code` to show a specific message rather than a generic fallback. Named SDK codes include `runtime_not_ready`, `asset_load_failed`, and -`worker_error`; daemon-originated errors currently use `wavelength_error`. +`worker_error`. When `start()` uses the default worker transport, it can raise +`runtime_locked` if another same-origin tab owns the runtime lock. +Daemon-originated errors currently use `wavelength_error`. ```tsx title="errors.tsx" import { useWalletSend, WavelengthError } from '@lightninglabs/wavelength-react'; diff --git a/apps/docs/src/content/docs/reference/wavelength-core.mdx b/apps/docs/src/content/docs/reference/wavelength-core.mdx index 96f56f5b..4290a6cf 100644 --- a/apps/docs/src/content/docs/reference/wavelength-core.mdx +++ b/apps/docs/src/content/docs/reference/wavelength-core.mdx @@ -2794,6 +2794,7 @@ export const wavelengthErrorSig = `class WavelengthError extends Error { export const errorCodeSig = `type WavelengthErrorCode = | 'wavelength_error' | 'runtime_not_ready' + | 'runtime_locked' | 'asset_load_failed' | 'worker_error' | 'unsupported_facade_method' @@ -2829,6 +2830,7 @@ compatibility while still offering autocomplete on the known codes. |---|---| | `wavelength_error` | The generic default: any SDK error without a more specific code, and all daemon-originated failures today. | | `runtime_not_ready` | The wasm runtime is not callable: it exited before signaling ready, or its call entry point is missing once loading has finished. | +| `runtime_locked` | `start()` in the default worker mode cannot acquire the runtime lock because another same-origin tab owns it. | | `asset_load_failed` | `ready()` fails to load the runtime assets (wasm binary or its supporting files). | | `worker_error` | The worker transport's underlying Worker fails or crashes. | | `unsupported_facade_method` | A call names a method the daemon facade does not expose. | diff --git a/apps/web-wallet-demo/wavewalletdk-smoke.spec.js b/apps/web-wallet-demo/wavewalletdk-smoke.spec.js index 12abd643..9f5c2a30 100644 --- a/apps/web-wallet-demo/wavewalletdk-smoke.spec.js +++ b/apps/web-wallet-demo/wavewalletdk-smoke.spec.js @@ -36,6 +36,84 @@ async function createReadyWallet( await page.getByRole("button", { name: "I saved it" }).click(); } +test("a second tab retries after the owning runtime closes", async ({ + context, + page, +}, testInfo) => { + const baseURL = testInfo.project.use.baseURL; + const dataDir = `/wavewalletdk-smoke-lock-${Date.now()}`; + const swapDatabaseFileName = `/wavewalletdk-swaps-lock-${Date.now()}.db`; + const walletEntry = (targetPage) => + targetPage.getByRole("heading", { + name: /^(Create wallet|Unlock wallet)$/, + }); + + await page.goto("/"); + const firstStart = page.getByRole("button", { name: "Start runtime" }); + await expect(firstStart).toBeVisible({ timeout: 30000 }); + await configureRuntime(page, baseURL, dataDir, swapDatabaseFileName); + await firstStart.click(); + await expect(walletEntry(page)).toBeVisible({ timeout: 60000 }); + + const secondPage = await context.newPage(); + const secondPageMessages = []; + const recordSecondPageMessage = (line) => { + secondPageMessages.push(line); + if (process.env.WAVELENGTH_SMOKE_VERBOSE) { + console.log(line); + } + }; + secondPage.on("console", (message) => { + recordSecondPageMessage(`[${message.type()}] ${message.text()}`); + }); + secondPage.on("pageerror", (error) => { + recordSecondPageMessage(`[pageerror] ${error.message}`); + }); + const expectNoRuntimeFallback = () => { + const output = secondPageMessages.join("\n"); + for (const marker of [ + "falling back to :memory:", + "OPFS VFS unavailable", + "SQLITE_CANTOPEN", + "apply sqlite migrations", + ]) { + expect(output).not.toContain(marker); + } + }; + + await secondPage.goto("/"); + const secondStart = secondPage.getByRole("button", { + name: "Start runtime", + }); + await expect(secondStart).toBeVisible({ timeout: 30000 }); + await configureRuntime( + secondPage, + baseURL, + dataDir, + swapDatabaseFileName, + ); + await secondStart.click(); + + await expect( + secondPage.getByRole("heading", { name: "Runtime error" }), + ).toBeVisible(); + await expect( + secondPage.getByText( + "This wallet is already open in another tab. Close the other tab and try again.", + ), + ).toBeVisible(); + await expect(secondPage.locator("body")).not.toContainText("SQLITE_CANTOPEN"); + await expect(secondPage.locator("body")).not.toContainText( + "apply sqlite migrations", + ); + expectNoRuntimeFallback(); + + await page.close(); + await secondPage.getByRole("button", { name: "Try again" }).click(); + await expect(walletEntry(secondPage)).toBeVisible({ timeout: 60000 }); + expectNoRuntimeFallback(); +}); + test("wallet create and address state persist with OPFS SQLite", async ({ page, }, testInfo) => { diff --git a/packages/core/src/errors.ts b/packages/core/src/errors.ts index d4a4a48f..4e04ba89 100644 --- a/packages/core/src/errors.ts +++ b/packages/core/src/errors.ts @@ -9,6 +9,7 @@ export type WavelengthErrorCode = | 'wavelength_error' | 'runtime_not_ready' + | 'runtime_locked' | 'asset_load_failed' | 'worker_error' | 'unsupported_facade_method' diff --git a/packages/web/src/clients/runtime-lock.ts b/packages/web/src/clients/runtime-lock.ts new file mode 100644 index 00000000..15e56c30 --- /dev/null +++ b/packages/web/src/clients/runtime-lock.ts @@ -0,0 +1,126 @@ +import { WavelengthError } from '@lightninglabs/wavelength-core'; + +const RUNTIME_LOCK_NAME = 'lightninglabs:wavelength:worker-runtime'; +const RUNTIME_LOCKED_MESSAGE = + 'This wallet is already open in another tab. Close the other tab and try again.'; + +type RuntimeLockLease = { + release: () => void; +}; + +function clientDisposedError(): WavelengthError { + return new WavelengthError('Wavelength client disposed', 'worker_error'); +} + +/** + * Holds the worker runtime's origin-scoped Web Lock until its storage-owning + * lifetime ends. The lock is intentionally not keyed by dataDir: the daemon + * can open independently configured paths such as the swap database, plus + * paths selected by daemon defaults. + */ +export class WorkerRuntimeLock { + private lease: RuntimeLockLease | null = null; + private priorRelease: Promise = Promise.resolve(); + private rejectAcquisition: ((reason: unknown) => void) | null = null; + private disposed = false; + + acquire(): Promise { + if (this.disposed) { + return Promise.reject(clientDisposedError()); + } + if (this.lease) { + return Promise.resolve(false); + } + + return this.acquireAfterPriorRelease(); + } + + release(): void { + const lease = this.lease; + this.lease = null; + lease?.release(); + } + + dispose(): void { + this.disposed = true; + this.rejectAcquisition?.(clientDisposedError()); + this.release(); + } + + private async acquireAfterPriorRelease(): Promise { + await this.priorRelease; + if (this.disposed) { + throw clientDisposedError(); + } + if (this.lease) { + return false; + } + + const lockManager = globalThis.navigator?.locks; + if (!lockManager) { + return false; + } + + let releaseLease!: () => void; + const holdLease = new Promise((resolve) => { + releaseLease = resolve; + }); + let resolveAvailability!: (available: boolean) => void; + let rejectAvailability!: (reason: unknown) => void; + const availability = new Promise((resolve, reject) => { + resolveAvailability = resolve; + rejectAvailability = reject; + }); + this.rejectAcquisition = rejectAvailability; + + let requestCompletion: Promise; + try { + requestCompletion = lockManager.request( + RUNTIME_LOCK_NAME, + { mode: 'exclusive', ifAvailable: true }, + async (lock) => { + if (!lock) { + resolveAvailability(false); + + return; + } + if (this.disposed) { + rejectAvailability(clientDisposedError()); + + return; + } + this.lease = { release: releaseLease }; + resolveAvailability(true); + await holdLease; + }, + ); + } catch (err) { + this.rejectAcquisition = null; + throw err; + } + + this.priorRelease = requestCompletion.then( + () => undefined, + () => undefined, + ); + void requestCompletion.catch(rejectAvailability); + + let available: boolean; + try { + available = await availability; + } finally { + if (this.rejectAcquisition === rejectAvailability) { + this.rejectAcquisition = null; + } + } + if (!available) { + throw new WavelengthError(RUNTIME_LOCKED_MESSAGE, 'runtime_locked'); + } + if (this.disposed) { + this.release(); + throw clientDisposedError(); + } + + return true; + } +} diff --git a/packages/web/src/clients/worker-runtime-lock.test.ts b/packages/web/src/clients/worker-runtime-lock.test.ts new file mode 100644 index 00000000..1efddb1a --- /dev/null +++ b/packages/web/src/clients/worker-runtime-lock.test.ts @@ -0,0 +1,621 @@ +import assert from 'node:assert/strict'; +import { register } from 'node:module'; +import { describe, it } from 'node:test'; + +const typescriptURL = new URL( + '../../../../node_modules/typescript/lib/typescript.js', + import.meta.url, +).href; + +register( + `data:text/javascript,${encodeURIComponent(` + import * as ts from ${JSON.stringify(typescriptURL)}; + + export async function resolve(specifier, context, nextResolve) { + try { + return await nextResolve(specifier, context); + } catch (error) { + if (error?.code === 'ERR_MODULE_NOT_FOUND' && specifier.startsWith('.') && !specifier.endsWith('.ts')) { + return nextResolve(specifier + '.ts', context); + } + throw error; + } + } + + export async function load(url, context, nextLoad) { + const result = await nextLoad(url, context); + if (!url.endsWith('.ts')) return result; + return { + ...result, + source: ts.transpileModule(String(result.source), { + compilerOptions: { + module: ts.ModuleKind.ESNext, + target: ts.ScriptTarget.ES2022, + }, + }).outputText, + }; + } + `)}`, + import.meta.url, +); + +type WorkerMessage = { + id?: number; + method?: string; + params?: unknown; + $init?: unknown; +}; + +type RuntimeWorkerRequest = WorkerMessage & { + id: number; + method: string; +}; + +type RuntimeClient = { + ready(): Promise; + start(config: { network?: 'signet' }): Promise; + isRunning(): Promise; + callFacade( + method: 'start' | 'stop', + params?: unknown, + ): Promise; + subscribe(listener: (event: unknown) => void): () => void; + dispose(): void; +}; + +function deferred(): { promise: Promise; resolve: (value: T) => void } { + let resolve!: (value: T) => void; + const promise = new Promise((res) => { + resolve = res; + }); + + return { promise, resolve }; +} + +class RuntimeFakeWorker { + static instances: RuntimeFakeWorker[] = []; + readonly messages: WorkerMessage[] = []; + readonly responded = new Set(); + onmessage: ((event: MessageEvent) => void) | null = null; + onerror: ((event: ErrorEvent) => void) | null = null; + terminated = false; + + constructor(_url: string | URL) { + RuntimeFakeWorker.instances.push(this); + } + + postMessage(message: WorkerMessage): void { + this.messages.push(message); + } + + terminate(): void { + this.terminated = true; + } + + lifecycleRequests(): RuntimeWorkerRequest[] { + return this.messages.filter( + (message): message is RuntimeWorkerRequest => + typeof message.id === 'number' && + (message.method === 'start' || message.method === 'stop'), + ); + } + + request(method: string): RuntimeWorkerRequest { + const request = this.messages.find( + (message): message is RuntimeWorkerRequest => + message.method === method && + typeof message.id === 'number' && + !this.responded.has(message.id), + ); + assert.ok(request, `expected an unsettled ${method} worker request`); + + return request; + } + + resolve(method: string, result?: unknown): void { + const request = this.request(method); + this.responded.add(request.id); + this.onmessage?.({ + data: { id: request.id, ok: true, result }, + } as MessageEvent); + } + + reject(method: string, error: string): void { + const request = this.request(method); + this.responded.add(request.id); + this.onmessage?.({ + data: { id: request.id, ok: false, error }, + } as MessageEvent); + } + + fatal(message: string): void { + this.onmessage?.({ data: { fatal: { message } } } as MessageEvent); + } +} + +type LockRequest = { + name: string; + options: LockOptions; +}; + +class FakeLockManager { + readonly requests: LockRequest[] = []; + private held = false; + private acquisitionGate: Promise = Promise.resolve(); + + get locked(): boolean { + return this.held; + } + + deferAcquisition(): () => void { + const gate = deferred(); + this.acquisitionGate = gate.promise; + + return () => gate.resolve(); + } + + request( + name: string, + options: LockOptions, + callback: (lock: Lock | null) => T | PromiseLike, + ): Promise { + this.requests.push({ name, options }); + + return (async () => { + await this.acquisitionGate; + if (this.held) { + return callback(null); + } + + this.held = true; + try { + return await callback({ name, mode: 'exclusive' } as Lock); + } finally { + this.held = false; + } + })(); + } +} + +type RuntimeHarness = { + createClient(): RuntimeClient; + disposeClient(client: RuntimeClient): void; +}; + +function restoreGlobal( + property: 'Worker' | 'navigator', + descriptor: PropertyDescriptor | undefined, +): void { + if (descriptor) { + Object.defineProperty(globalThis, property, descriptor); + } else { + Reflect.deleteProperty(globalThis, property); + } +} + +async function flushMicrotasks(times = 12): Promise { + for (let index = 0; index < times; index += 1) { + await Promise.resolve(); + } +} + +async function withRuntimeHarness( + lockManager: FakeLockManager | null, + run: (harness: RuntimeHarness) => Promise, +): Promise { + const savedWorker = Object.getOwnPropertyDescriptor(globalThis, 'Worker'); + const savedNavigator = Object.getOwnPropertyDescriptor( + globalThis, + 'navigator', + ); + const clients = new Set(); + RuntimeFakeWorker.instances = []; + Object.defineProperty(globalThis, 'Worker', { + configurable: true, + value: RuntimeFakeWorker, + }); + Object.defineProperty(globalThis, 'navigator', { + configurable: true, + value: lockManager ? { locks: lockManager } : {}, + }); + + try { + const { WorkerWavelengthClient } = await import('./worker.ts'); + const createClient = (): RuntimeClient => { + const client = new WorkerWavelengthClient({ + workerURL: 'fake-worker.js', + }) as RuntimeClient; + clients.add(client); + + return client; + }; + await run({ + createClient, + disposeClient: (client) => { + client.dispose(); + clients.delete(client); + }, + }); + } finally { + for (const client of clients) { + client.dispose(); + } + restoreGlobal('Worker', savedWorker); + restoreGlobal('navigator', savedNavigator); + RuntimeFakeWorker.instances = []; + } +} + +describe('worker runtime lock', () => { + it('rejects a second runtime before dispatching start', async () => { + const lockManager = new FakeLockManager(); + await withRuntimeHarness(lockManager, async ({ createClient }) => { + const { WavelengthError } = await import( + '@lightninglabs/wavelength-core' + ); + const firstClient = createClient(); + const secondClient = createClient(); + const firstWorker = RuntimeFakeWorker.instances[0]; + const secondWorker = RuntimeFakeWorker.instances[1]; + + const firstStart = firstClient.callFacade('start'); + await flushMicrotasks(); + assert.equal(firstWorker.request('start').method, 'start'); + assert.deepEqual(lockManager.requests[0], { + name: 'lightninglabs:wavelength:worker-runtime', + options: { mode: 'exclusive', ifAvailable: true }, + }); + firstWorker.resolve('start'); + await firstStart; + + await assert.rejects( + () => secondClient.start({ network: 'signet' }), + (err: unknown) => { + assert.ok(err instanceof WavelengthError); + assert.equal(err.code, 'runtime_locked'); + assert.equal( + err.message, + 'This wallet is already open in another tab. Close the other tab and try again.', + ); + + return true; + }, + ); + assert.equal(secondWorker.lifecycleRequests().length, 0); + }); + }); + + it('continues queued lifecycle work after a rejected start', async () => { + const lockManager = new FakeLockManager(); + await withRuntimeHarness(lockManager, async ({ createClient }) => { + const client = createClient(); + const worker = RuntimeFakeWorker.instances[0]; + const firstStart = client.callFacade('start', { marker: 'first' }); + const secondStart = client.callFacade('start', { marker: 'second' }); + + await flushMicrotasks(); + assert.deepEqual( + worker.lifecycleRequests().map(({ method }) => method), + ['start'], + ); + worker.reject('start', 'start failed'); + await assert.rejects(() => firstStart, /start failed/); + await flushMicrotasks(); + + assert.deepEqual( + worker.lifecycleRequests().map(({ method }) => method), + ['start', 'start'], + ); + assert.deepEqual(worker.lifecycleRequests()[1].params, { + marker: 'second', + }); + worker.resolve('start', 'second result'); + assert.equal(await secondStart, 'second result'); + assert.equal(lockManager.requests.length, 2); + assert.equal(lockManager.locked, true); + }); + }); + + it('keeps an existing lock when a repeated start fails', async () => { + const lockManager = new FakeLockManager(); + await withRuntimeHarness(lockManager, async ({ createClient }) => { + const firstClient = createClient(); + const firstWorker = RuntimeFakeWorker.instances[0]; + const firstStart = firstClient.callFacade('start'); + await flushMicrotasks(); + firstWorker.resolve('start'); + await firstStart; + + const repeatedStart = firstClient.callFacade('start'); + await flushMicrotasks(); + firstWorker.reject('start', 'already running'); + await assert.rejects(() => repeatedStart, /already running/); + + const secondClient = createClient(); + const secondWorker = RuntimeFakeWorker.instances[1]; + await assert.rejects( + () => secondClient.callFacade('start'), + (err: { code?: string }) => err.code === 'runtime_locked', + ); + assert.equal(secondWorker.lifecycleRequests().length, 0); + }); + }); + + it('keeps the lock after a failed stop', async () => { + const lockManager = new FakeLockManager(); + await withRuntimeHarness(lockManager, async ({ createClient }) => { + const firstClient = createClient(); + const firstWorker = RuntimeFakeWorker.instances[0]; + const firstStart = firstClient.callFacade('start'); + await flushMicrotasks(); + firstWorker.resolve('start'); + await firstStart; + + const stop = firstClient.callFacade('stop'); + await flushMicrotasks(); + firstWorker.reject('stop', 'stop failed'); + await assert.rejects(() => stop, /stop failed/); + + const secondClient = createClient(); + const secondWorker = RuntimeFakeWorker.instances[1]; + await assert.rejects( + () => secondClient.callFacade('start'), + (err: { code?: string }) => err.code === 'runtime_locked', + ); + assert.equal(secondWorker.lifecycleRequests().length, 0); + }); + }); + + it('orders pending start, stop, and start calls exactly', async () => { + const lockManager = new FakeLockManager(); + await withRuntimeHarness(lockManager, async ({ createClient }) => { + const firstClient = createClient(); + const firstWorker = RuntimeFakeWorker.instances[0]; + const startA = firstClient.callFacade('start', { marker: 'start A' }); + const stopB = firstClient.callFacade('stop', { marker: 'stop B' }); + const startC = firstClient.callFacade('start', { marker: 'start C' }); + assert.notEqual(startA, startC); + + await flushMicrotasks(); + assert.deepEqual( + firstWorker.lifecycleRequests().map(({ method }) => method), + ['start'], + ); + assert.deepEqual(firstWorker.lifecycleRequests()[0].params, { + marker: 'start A', + }); + + firstWorker.resolve('start', 'result A'); + assert.equal(await startA, 'result A'); + await flushMicrotasks(); + assert.deepEqual( + firstWorker.lifecycleRequests().map(({ method }) => method), + ['start', 'stop'], + ); + + firstWorker.resolve('stop', 'result B'); + assert.equal(await stopB, 'result B'); + await flushMicrotasks(); + assert.deepEqual( + firstWorker.lifecycleRequests().map(({ method }) => method), + ['start', 'stop', 'start'], + ); + assert.deepEqual(firstWorker.lifecycleRequests()[2].params, { + marker: 'start C', + }); + + firstWorker.resolve('start', 'result C'); + assert.equal(await startC, 'result C'); + assert.equal(lockManager.requests.length, 2); + + const secondClient = createClient(); + const secondWorker = RuntimeFakeWorker.instances[1]; + await assert.rejects( + () => secondClient.callFacade('start'), + (err: { code?: string }) => err.code === 'runtime_locked', + ); + assert.equal(secondWorker.lifecycleRequests().length, 0); + }); + }); + + it('orders pending stop, start, and stop calls exactly', async () => { + const lockManager = new FakeLockManager(); + await withRuntimeHarness(lockManager, async ({ createClient }) => { + const firstClient = createClient(); + const firstWorker = RuntimeFakeWorker.instances[0]; + const initialStart = firstClient.callFacade('start'); + await flushMicrotasks(); + firstWorker.resolve('start'); + await initialStart; + + const requestOffset = firstWorker.lifecycleRequests().length; + const stopA = firstClient.callFacade('stop', { marker: 'stop A' }); + const startB = firstClient.callFacade('start', { marker: 'start B' }); + const stopC = firstClient.callFacade('stop', { marker: 'stop C' }); + assert.notEqual(stopA, stopC); + + await flushMicrotasks(); + assert.deepEqual( + firstWorker + .lifecycleRequests() + .slice(requestOffset) + .map(({ method }) => method), + ['stop'], + ); + + firstWorker.resolve('stop', 'result A'); + assert.equal(await stopA, 'result A'); + await flushMicrotasks(); + assert.deepEqual( + firstWorker + .lifecycleRequests() + .slice(requestOffset) + .map(({ method }) => method), + ['stop', 'start'], + ); + + firstWorker.resolve('start', 'result B'); + assert.equal(await startB, 'result B'); + await flushMicrotasks(); + assert.deepEqual( + firstWorker + .lifecycleRequests() + .slice(requestOffset) + .map(({ method }) => method), + ['stop', 'start', 'stop'], + ); + + firstWorker.resolve('stop', 'result C'); + assert.equal(await stopC, 'result C'); + assert.equal(lockManager.locked, false); + + const secondClient = createClient(); + const secondWorker = RuntimeFakeWorker.instances[1]; + const secondStart = secondClient.callFacade('start'); + await flushMicrotasks(); + assert.equal(secondWorker.request('start').method, 'start'); + secondWorker.resolve('start'); + await secondStart; + }); + }); + + it('releases the lock and terminates the worker on dispose', async () => { + const lockManager = new FakeLockManager(); + await withRuntimeHarness( + lockManager, + async ({ createClient, disposeClient }) => { + const firstClient = createClient(); + const firstWorker = RuntimeFakeWorker.instances[0]; + const firstStart = firstClient.callFacade('start'); + await flushMicrotasks(); + firstWorker.resolve('start'); + await firstStart; + + disposeClient(firstClient); + assert.equal(firstWorker.terminated, true); + await flushMicrotasks(); + assert.equal(lockManager.locked, false); + + const secondClient = createClient(); + const secondWorker = RuntimeFakeWorker.instances[1]; + const secondStart = secondClient.callFacade('start'); + await flushMicrotasks(); + assert.equal(secondWorker.request('start').method, 'start'); + secondWorker.resolve('start'); + await secondStart; + }, + ); + }); + + it('disposes safely while Web Lock acquisition is pending', async () => { + const lockManager = new FakeLockManager(); + const continueAcquisition = lockManager.deferAcquisition(); + await withRuntimeHarness( + lockManager, + async ({ createClient, disposeClient }) => { + const firstClient = createClient(); + const firstWorker = RuntimeFakeWorker.instances[0]; + const firstStart = firstClient.callFacade('start'); + await flushMicrotasks(); + assert.equal(lockManager.requests.length, 1); + assert.equal(firstWorker.lifecycleRequests().length, 0); + + disposeClient(firstClient); + assert.equal(firstWorker.terminated, true); + await assert.rejects( + () => firstStart, + (err: { code?: string; message?: string }) => + err.code === 'worker_error' && + err.message === 'Wavelength client disposed', + ); + + continueAcquisition(); + await flushMicrotasks(); + assert.equal(firstWorker.lifecycleRequests().length, 0); + assert.equal(lockManager.locked, false); + + const secondClient = createClient(); + const secondWorker = RuntimeFakeWorker.instances[1]; + const secondStart = secondClient.callFacade('start'); + await flushMicrotasks(); + assert.equal(secondWorker.request('start').method, 'start'); + secondWorker.resolve('start'); + await secondStart; + }, + ); + }); + + it('releases the lock after a fatal runtime exit', async () => { + const lockManager = new FakeLockManager(); + await withRuntimeHarness(lockManager, async ({ createClient }) => { + const firstClient = createClient(); + const firstWorker = RuntimeFakeWorker.instances[0]; + const events: unknown[] = []; + firstClient.subscribe((event) => events.push(event)); + const firstStart = firstClient.callFacade('start'); + await flushMicrotasks(); + firstWorker.resolve('start'); + await firstStart; + + firstWorker.fatal('runtime exited'); + const secondClient = createClient(); + const secondWorker = RuntimeFakeWorker.instances[1]; + const secondStart = secondClient.callFacade('start'); + await flushMicrotasks(); + assert.deepEqual(events.at(-1), { type: 'runtimeStopped' }); + assert.equal(secondWorker.request('start').method, 'start'); + secondWorker.resolve('start'); + await secondStart; + }); + }); + + it('keeps the lock when post-start getInfo fails', async () => { + const lockManager = new FakeLockManager(); + await withRuntimeHarness(lockManager, async ({ createClient }) => { + const firstClient = createClient(); + const firstWorker = RuntimeFakeWorker.instances[0]; + const start = firstClient.start({ network: 'signet' }); + await flushMicrotasks(); + firstWorker.resolve('start'); + await flushMicrotasks(); + firstWorker.reject('getInfo', 'getInfo failed'); + await assert.rejects(() => start, /getInfo failed/); + + const secondClient = createClient(); + const secondWorker = RuntimeFakeWorker.instances[1]; + await assert.rejects( + () => secondClient.callFacade('start'), + (err: { code?: string }) => err.code === 'runtime_locked', + ); + assert.equal(secondWorker.lifecycleRequests().length, 0); + }); + }); + + it('does not lock ready or unrelated facade requests', async () => { + const lockManager = new FakeLockManager(); + await withRuntimeHarness(lockManager, async ({ createClient }) => { + const client = createClient(); + const worker = RuntimeFakeWorker.instances[0]; + + const ready = client.ready(); + worker.resolve('$ready'); + await ready; + const isRunning = client.isRunning(); + worker.resolve('isRunning', true); + assert.equal(await isRunning, true); + assert.equal(lockManager.requests.length, 0); + }); + }); + + it('preserves worker behavior when Web Locks are unavailable', async () => { + await withRuntimeHarness(null, async ({ createClient }) => { + const client = createClient(); + const worker = RuntimeFakeWorker.instances[0]; + const start = client.callFacade('start'); + await flushMicrotasks(); + assert.equal(worker.request('start').method, 'start'); + worker.resolve('start'); + await start; + }); + }); +}); diff --git a/packages/web/src/clients/worker.ts b/packages/web/src/clients/worker.ts index e5c19772..c388126d 100644 --- a/packages/web/src/clients/worker.ts +++ b/packages/web/src/clients/worker.ts @@ -10,6 +10,7 @@ import type { import type { WebClientOptions } from '../index.ts'; import { defaultWorkerRuntimeBaseUrl } from '../runtime.ts'; import { PendingCall, errorMessage, toWavelengthEvent } from '../util.ts'; +import { WorkerRuntimeLock } from './runtime-lock.ts'; type WorkerControlMethod = '$ready' | '$startActivity' | '$stopActivity'; @@ -25,7 +26,10 @@ export class WorkerWavelengthClient extends BaseWavelengthClient { protected readonly serverTransport = 'rest' as const; private readonly worker: Worker; private readonly pending = new Map(); + private readonly runtimeLock = new WorkerRuntimeLock(); + private lifecycleTail: Promise = Promise.resolve(); private nextRequestID = 1; + private disposed = false; constructor(options: WebClientOptions = {}) { super(); @@ -63,9 +67,70 @@ export class WorkerWavelengthClient extends BaseWavelengthClient { method: FacadeMethod, params: unknown = {}, ): Promise { + if (method === 'start') { + return this.startRuntime(params) as Promise; + } + if (method === 'stop') { + return this.stopRuntime(params) as Promise; + } + return this.request(method, params); } + private startRuntime(params: unknown): Promise { + return this.enqueueLifecycle(() => this.dispatchStart(params)); + } + + private async dispatchStart(params: unknown): Promise { + this.assertNotDisposed(); + const lockAcquiredForStart = await this.runtimeLock.acquire(); + try { + this.assertNotDisposed(); + + return await this.request('start', params); + } catch (err) { + // A rejected start facade request means the daemon never completed + // startup. Release only when this call acquired the lock: a repeated + // start against an already-running daemon must not unlock its existing + // storage ownership. This catch runs before BaseWavelengthClient's + // separate post-start getInfo request. + if (lockAcquiredForStart) { + this.runtimeLock.release(); + } + throw err; + } + } + + private stopRuntime(params: unknown): Promise { + return this.enqueueLifecycle(() => this.dispatchStop(params)); + } + + private async dispatchStop(params: unknown): Promise { + this.assertNotDisposed(); + const result = await this.request('stop', params); + // A failed stop may leave the daemon's OPFS handles alive, so release only + // after the facade confirms the runtime stopped. + this.runtimeLock.release(); + + return result; + } + + private enqueueLifecycle(operation: () => Promise): Promise { + const request = this.lifecycleTail.then(operation, operation); + this.lifecycleTail = request.then( + () => undefined, + () => undefined, + ); + + return request; + } + + private assertNotDisposed(): void { + if (this.disposed) { + throw new WavelengthError('Wavelength client disposed', 'worker_error'); + } + } + private request( method: FacadeMethod | WorkerControlMethod, params: unknown = {}, @@ -130,6 +195,7 @@ export class WorkerWavelengthClient extends BaseWavelengthClient { // not hang, and emit runtimeStopped so subscribers (e.g. the provider) // move the lifecycle off 'ready' instead of appearing alive after the // engine died. + this.runtimeLock.release(); this.rejectAll( new WavelengthError( data.fatal.message || 'Wavelength worker stopped', @@ -184,12 +250,14 @@ export class WorkerWavelengthClient extends BaseWavelengthClient { } dispose(): void { + this.disposed = true; super.dispose(); // Terminating the worker fires neither onerror nor a fatal message, so // reject any in-flight calls here so they do not hang past disposal. this.rejectAll( new WavelengthError('Wavelength client disposed', 'worker_error'), ); + this.runtimeLock.dispose(); this.worker.terminate(); } }