import type { Log, Message } from '../types.js'; import type Messenger from './Messenger.js'; import type { InitializeMessengerOptions, MessageHandler, } from './Messenger.js'; import { isAck2Message, isAck1Message, isSynMessage } from '../guards.js'; import PenpalError from '../PenpalError.js'; // Keep this structural so generated .d.ts files don't depend on // DedicatedWorkerGlobalScope being present in consumer tsconfig libs. type WorkerLike = Pick< Worker, 'postMessage' | 'addEventListener' | 'removeEventListener' >; type Options = { /** * The web worker receiving/sending communication from/to the parent window. * If this messenger is being used within the worker, `worker` should * typically be set to `self`. */ worker: WorkerLike; }; /** * Handles the details of communicating with a child web worker. */ class WorkerMessenger implements Messenger { #worker: WorkerLike; #log: Log | undefined; #validateReceivedMessage: ((data: unknown) => data is Message) | undefined; #messageCallbacks = new Set(); #port: MessagePort | undefined; constructor({ worker }: Options) { if (!worker) { throw new PenpalError('INVALID_ARGUMENT', 'worker must be defined'); } this.#worker = worker; } initialize = ({ log, validateReceivedMessage, }: InitializeMessengerOptions): void => { this.#log = log; this.#validateReceivedMessage = validateReceivedMessage; this.#worker.addEventListener('message', this.#handleMessage); }; sendMessage = (message: Message, transferables?: Transferable[]): void => { if (isSynMessage(message) || isAck1Message(message)) { this.#worker.postMessage(message, { ...(transferables === undefined ? {} : { transfer: transferables }), }); return; } if (isAck2Message(message)) { const { port1, port2 } = new MessageChannel(); this.#setPort(port1); try { this.#worker.postMessage(message, { transfer: [port2, ...(transferables || [])], }); } catch (error) { this.#destroyPort(); port2.close(); throw error; } return; } if (this.#port) { this.#port.postMessage(message, { ...(transferables === undefined ? {} : { transfer: transferables }), }); return; } throw new PenpalError( 'TRANSMISSION_FAILED', 'Cannot send message because the MessagePort is not connected', ); }; addMessageHandler = (callback: MessageHandler): void => { this.#messageCallbacks.add(callback); }; removeMessageHandler = (callback: MessageHandler): void => { this.#messageCallbacks.delete(callback); }; destroy = (): void => { this.#worker.removeEventListener('message', this.#handleMessage); this.#destroyPort(); this.#messageCallbacks.clear(); }; #destroyPort = () => { this.#port?.removeEventListener('message', this.#handleMessage); this.#port?.close(); this.#port = undefined; }; #setPort = (port: MessagePort) => { this.#destroyPort(); this.#port = port; this.#port.addEventListener('message', this.#handleMessage); this.#port.start(); }; #handleMessage = ({ ports, data }: MessageEvent): void => { if (!this.#validateReceivedMessage?.(data)) { return; } if (isSynMessage(data)) { // If we receive a SYN message and already have a port, it means // the child is re-connecting, in which case we'll receive a new port. // For this reason, we always make sure we destroy the existing port. this.#destroyPort(); } if (isAck2Message(data)) { const port = ports[0]; if (!port) { this.#log?.('Ignoring ACK2 because it did not include a MessagePort'); return; } this.#setPort(port); } for (const callback of this.#messageCallbacks) { callback(data); } }; } export default WorkerMessenger;