diff --git a/.changeset/lucky-otters-listen.md b/.changeset/lucky-otters-listen.md new file mode 100644 index 0000000..0e8efcc --- /dev/null +++ b/.changeset/lucky-otters-listen.md @@ -0,0 +1,13 @@ +--- +'webactor': minor +--- + +Notice a worker that dies before its handshake answers. + +A worker supervisor watches its worker's liveness through a lock whose key it learns from the handshake reply, so the watch can only be armed once that reply arrives. Before it did, the only detector was the worker's own `error` event. A worker that died quietly in that window — killed by the host, or closing itself — left the supervisor holding a dead worker forever: no restart, no error, nothing. A worker that came up but never answered was equally invisible. + +`applyWorkerSupervisor` now accepts `getAbortSignal`, a factory consulted once per launch, so the handshake can be bounded the same way every other operation in the library is: `getAbortSignal: () => AbortSignal.timeout(2000)`. It is a factory rather than a plain signal because a supervisor relaunches, and one signal would already be spent by the second worker. Whatever the returned signal aborts with reaches `shouldRetry` as the reason, so an `AbortSignal.timeout` arrives as a `TimeoutError`. + +Nothing changes when it is omitted. A deadline tight enough to be useful would misfire on a loaded machine, where a handshake legitimately takes hundreds of milliseconds, so the choice stays with the caller. + +A handshake that fails for any reason other than the supervisor's own teardown now reaches the restart decision too. It previously became an unhandled rejection instead, which meant an undeliverable handshake was reported to nobody and restarted nothing. diff --git a/packages/webactor/src/worker/applyWorkerSupervisor.ts b/packages/webactor/src/worker/applyWorkerSupervisor.ts index 6c06e37..38870d5 100644 --- a/packages/webactor/src/worker/applyWorkerSupervisor.ts +++ b/packages/webactor/src/worker/applyWorkerSupervisor.ts @@ -21,8 +21,10 @@ export function applyWorkerSupervisor( WorkerConstructor: () => Worker, { shouldRetry, + getAbortSignal, }: { shouldRetry: (reason?: unknown | Reason | Error | ErrorEvent) => boolean | Promise; + getAbortSignal?: () => undefined | AbortSignal; }, ): Actor { const proxy = createEnvelopeChannel(); @@ -52,13 +54,22 @@ export function applyWorkerSupervisor( .catch(catchAbortToSymbol); }; - request(messagePort, THREAD_ID_REQUEST, { abortSignal: abortController.signal }) + const callerSignal = getAbortSignal?.(); + const handshakeSignal = + callerSignal === undefined + ? abortController.signal + : AbortSignal.any([abortController.signal, callerSignal]); + + request(messagePort, THREAD_ID_REQUEST, { abortSignal: handshakeSignal }) .then((envelope) => { if (isObject(envelope.data) && isStringField(envelope.data, 'threadId')) { onUnlockThreadId(envelope.data.threadId); } }) - .catch(catchAbortToSymbol); + .catch((error) => { + if (abortController.signal.aborted) return; + decide(error); + }); const errorOff = on(worker, 'error', (error) => decide(error)); diff --git a/packages/webactor/tests/worker/supervisor.test.ts b/packages/webactor/tests/worker/supervisor.test.ts index 686244e..c92fe96 100644 --- a/packages/webactor/tests/worker/supervisor.test.ts +++ b/packages/webactor/tests/worker/supervisor.test.ts @@ -1,9 +1,9 @@ import '../locks'; import { Worker } from '@apacheli/web-workers'; -import { afterEach, describe, expect, it } from 'vitest'; +import { afterEach, describe, expect, it, vi } from 'vitest'; -import { Actor } from '../../src/types'; +import { Actor, AnyData } from '../../src/types'; import { applyWorkerSupervisor } from '../../src/worker/applyWorkerSupervisor'; function createWorker() { @@ -18,6 +18,10 @@ function createErrorWorker() { }); } +function isTaggedData(data: AnyData, type: string) { + return typeof data === 'object' && data !== null && 'type' in data && data.type === type; +} + describe('Worker Supervisor Tests with Real Workers', () => { let supervisedActor: Actor; let workers: Worker[] = []; @@ -123,7 +127,8 @@ describe('Worker Supervisor Tests with Real Workers', () => { }); supervisedActor.launch(); - await new Promise((resolve) => setTimeout(resolve, 100)); + await vi.waitFor(() => expect(retryCount).toBeGreaterThanOrEqual(1), { timeout: 5000, interval: 10 }); + await new Promise((resolve) => setTimeout(resolve, 200)); expect(retryCount).toBe(1); expect(createCount).toBe(1); @@ -223,15 +228,16 @@ describe('Worker Supervisor Tests with Real Workers', () => { supervisedActor.launch(); - // Wait for error worker to fail and restart cycles - await new Promise((resolve) => setTimeout(resolve, 100)); + await vi.waitFor( + () => { + expect(createCount).toBeGreaterThan(1); + expect(restartReasons.length).toBeGreaterThan(0); + expect(workers.length).toBeGreaterThan(1); + }, + { timeout: 5000, interval: 10 }, + ); console.log(`Final state: createCount=${createCount}, restartReasons:`, restartReasons); - - // Should have created multiple workers due to errors - expect(createCount).toBeGreaterThan(1); - expect(restartReasons.length).toBeGreaterThan(0); - expect(workers.length).toBeGreaterThan(1); }); it('should restart worker when terminated manually', async () => { @@ -243,15 +249,6 @@ describe('Worker Supervisor Tests with Real Workers', () => { console.log(`Creating terminate-test worker #${createCount}`); const worker = createWorker(); workers.push(worker); - - // Terminate the first worker after short delay - if (createCount === 1) { - setTimeout(() => { - console.log('Manually terminating first worker...'); - worker.terminate(); - }, 300); - } - return worker; }; @@ -267,16 +264,175 @@ describe('Worker Supervisor Tests with Real Workers', () => { }, }); + let pongs = 0; + supervisedActor.addEventListener('message', (envelope) => { + if (isTaggedData(envelope.data, 'pong')) pongs++; + }); + supervisedActor.launch(); - // Wait for termination and restart - await new Promise((resolve) => setTimeout(resolve, 300)); + // terminating before the worker has connected leaves nothing for the supervisor to notice + await vi.waitFor( + () => { + supervisedActor.postMessage({ type: 'ping' }); + expect(pongs).toBeGreaterThan(0); + }, + { timeout: 5000, interval: 20 }, + ); + + console.log('Manually terminating first worker...'); + workers[0].terminate(); + + await vi.waitFor( + () => { + expect(createCount).toBeGreaterThanOrEqual(2); + expect(workers.length).toBeGreaterThanOrEqual(2); + }, + { timeout: 5000, interval: 10 }, + ); console.log(`Termination test result: createCount=${createCount}, restartReasons:`, restartReasons); + }); + + it('should detect a worker that dies before the handshake when given an abort signal', async () => { + let createCount = 0; + const reasons: unknown[] = []; + + const workerConstructor = () => { + createCount++; + const worker = createWorker(); + workers.push(worker); + if (createCount === 1) worker.terminate(); + return worker; + }; + + supervisedActor = applyWorkerSupervisor(workerConstructor, { + getAbortSignal: () => AbortSignal.timeout(150), + shouldRetry: async (reason) => { + reasons.push(reason); + return false; + }, + }); + + supervisedActor.launch(); + + await vi.waitFor(() => expect(reasons.length).toBeGreaterThan(0), { timeout: 5000, interval: 10 }); + + expect(reasons[0]).toBeInstanceOf(Error); + expect(String(reasons[0])).toContain('TimeoutError'); + expect(createCount).toBe(1); + }); + + it('should leave a worker that dies before the handshake unnoticed without an abort signal', async () => { + let createCount = 0; + let decisions = 0; + + const workerConstructor = () => { + createCount++; + const worker = createWorker(); + workers.push(worker); + if (createCount === 1) worker.terminate(); + return worker; + }; + + supervisedActor = applyWorkerSupervisor(workerConstructor, { + shouldRetry: async () => { + decisions++; + return false; + }, + }); + + supervisedActor.launch(); + await new Promise((resolve) => setTimeout(resolve, 600)); + + expect(decisions).toBe(0); + expect(createCount).toBe(1); + }); + + it('should build a fresh abort signal for every relaunch', async () => { + let createCount = 0; + let signalCount = 0; + const reasons: unknown[] = []; + + const workerConstructor = () => { + createCount++; + const worker = createWorker(); + workers.push(worker); + if (createCount <= 2) worker.terminate(); + return worker; + }; + + supervisedActor = applyWorkerSupervisor(workerConstructor, { + getAbortSignal: () => { + signalCount++; + return AbortSignal.timeout(150); + }, + shouldRetry: async (reason) => { + reasons.push(reason); + return reasons.length < 2; + }, + }); + + supervisedActor.launch(); + + // a single signal would already be spent here, so the second worker would go unwatched + await vi.waitFor(() => expect(reasons.length).toBe(2), { timeout: 5000, interval: 10 }); + + expect(signalCount).toBe(createCount); + expect(createCount).toBe(2); + }); + + it('should treat a plain abort from the caller as a failed handshake', async () => { + let createCount = 0; + const reasons: unknown[] = []; + const abortController = new AbortController(); + + const workerConstructor = () => { + createCount++; + const worker = createWorker(); + workers.push(worker); + if (createCount === 1) worker.terminate(); + return worker; + }; + + supervisedActor = applyWorkerSupervisor(workerConstructor, { + getAbortSignal: () => abortController.signal, + shouldRetry: async (reason) => { + reasons.push(reason); + return false; + }, + }); + + supervisedActor.launch(); + abortController.abort(); + + await vi.waitFor(() => expect(reasons.length).toBeGreaterThan(0), { timeout: 5000, interval: 10 }); + + expect(createCount).toBe(1); + }); + + it('should not ask for a restart decision when the supervisor itself is closed', async () => { + let decisions = 0; + + const workerConstructor = () => { + const worker = createWorker(); + workers.push(worker); + return worker; + }; + + supervisedActor = applyWorkerSupervisor(workerConstructor, { + getAbortSignal: () => AbortSignal.timeout(5000), + shouldRetry: async () => { + decisions++; + return false; + }, + }); + + supervisedActor.launch(); + supervisedActor.close(); + await new Promise((resolve) => setTimeout(resolve, 300)); - // Should have restarted after termination - expect(createCount).toBeGreaterThanOrEqual(2); - expect(workers.length).toBeGreaterThanOrEqual(2); + expect(decisions).toBe(0); }); it('should handle worker that throws error on message', async () => {