import {Actor, type ActorTarget, type MessageHandler} from './actor.ts'; import {getGlobalWorkerPool} from './global_worker_pool.ts'; import {GLOBAL_DISPATCHER_ID, makeRequest} from './ajax.ts'; import type {WorkerPool} from './worker_pool.ts'; import type {RequestResponseMessageMap} from './actor_messages.ts'; import {MessageType} from './actor_messages.ts'; /** * Responsible for sending messages from a {@link Source} to an associated worker source (usually with the same name). */ export class Dispatcher { workerPool: WorkerPool; actors: Actor[]; actorsPromise: Promise; currentActor: number; id: string | number; private removed: boolean; constructor(workerPool: WorkerPool, mapId: string | number) { this.workerPool = workerPool; this.actors = []; this.currentActor = 0; this.id = mapId; this.removed = false; this.actorsPromise = this.initActors(mapId); } private async initActors(mapId: string | number): Promise { const workers = await this.workerPool.acquire(mapId); if (this.removed) return []; this.actors = workers.map((worker: ActorTarget, i: number) => { const actor = new Actor(worker, mapId); actor.name = `Worker ${i}`; return actor; }); if (!this.actors.length) throw new Error('No actors found'); return this.actors; } /** * Broadcast a message to all Workers. */ async broadcast(type: T, data: RequestResponseMessageMap[T][0]): Promise> { const actors = await this.actorsPromise; return Promise.all(actors.map(actor => actor.sendAsync({type, data}))); } /** * Acquires an actor to dispatch messages to. The actors are distributed in round-robin fashion. * @returns An actor object backed by a web worker for processing messages. */ async getActor(): Promise { const actors = await this.actorsPromise; this.currentActor = (this.currentActor + 1) % actors.length; return actors[this.currentActor]; } async waitForInitComplete(): Promise { if (this.actors.length === 0) { await this.actorsPromise; } } getReadyActor(): Actor { this.currentActor = (this.currentActor + 1) % this.actors.length; return this.actors[this.currentActor]; } remove(mapRemoved: boolean = true): void { this.removed = true; for (const actor of this.actors) { actor.remove(); } this.actors = []; if (mapRemoved) this.workerPool.release(this.id); } public async registerMessageHandler(type: T, handler: MessageHandler): Promise { const actors = await this.actorsPromise; for (const actor of actors) { actor.registerMessageHandler(type, handler); } } public async unregisterMessageHandler(type: T): Promise { const actors = await this.actorsPromise; for (const actor of actors) { actor.unregisterMessageHandler(type); } } } let globalDispatcher: Dispatcher; /** * This function is used to get the global dispatcher that is shared across all maps instances. * It is used by the main thread to send messages to the workers, and by the workers to send messages back to the main thread. * If you import a script into the worker and need to send a message to the workers to pass some parameters for example, * you can use this function to get the global dispatcher and send a message to the workers. * @returns The global dispatcher instance. */ export function getGlobalDispatcher(): Dispatcher { if (!globalDispatcher) { globalDispatcher = new Dispatcher(getGlobalWorkerPool(), GLOBAL_DISPATCHER_ID); globalDispatcher.registerMessageHandler(MessageType.getResource, (_mapId, params, abortController) => { return makeRequest(params, abortController); }); } return globalDispatcher; }