this.opts.workerChoiceStrategyOptions =
opts.workerChoiceStrategyOptions ??
DEFAULT_WORKER_CHOICE_STRATEGY_OPTIONS
+ this.checkValidWorkerChoiceStrategyOptions(
+ this.opts.workerChoiceStrategyOptions
+ )
this.opts.enableEvents = opts.enableEvents ?? true
this.opts.enableTasksQueue = opts.enableTasksQueue ?? false
if (this.opts.enableTasksQueue) {
'Invalid worker choice strategy options: must be a plain object'
)
}
+ if (
+ workerChoiceStrategyOptions.weights != null &&
+ Object.keys(workerChoiceStrategyOptions.weights).length !== this.size
+ ) {
+ throw new Error(
+ 'Invalid worker choice strategy options: must have a weight for each worker node'
+ )
+ }
}
private checkValidTasksQueueOptions (
protected internalBusy (): boolean {
return (
this.workerNodes.findIndex(workerNode => {
- return workerNode.tasksUsage?.running === 0
+ return workerNode.tasksUsage.running === 0
}) === -1
)
}
/** @inheritDoc */
public async execute (data?: Data, name?: string): Promise<Response> {
- const [workerNodeKey, workerNode] = this.chooseWorkerNode()
+ const workerNodeKey = this.chooseWorkerNode()
const submittedTask: Task<Data> = {
name,
// eslint-disable-next-line @typescript-eslint/consistent-type-assertions
this.promiseResponseMap.set(submittedTask.id as string, {
resolve,
reject,
- worker: workerNode.worker
+ worker: this.workerNodes[workerNodeKey].worker
})
})
if (
} else {
this.executeTask(workerNodeKey, submittedTask)
}
+ this.workerChoiceStrategyContext.update(workerNodeKey)
this.checkAndEmitEvents()
// eslint-disable-next-line @typescript-eslint/return-await
return res
worker: Worker,
message: MessageValue<Response>
): void {
- const workerTasksUsage = this.getWorkerTasksUsage(worker)
+ const workerNodeKey = this.getWorkerNodeKey(worker)
+ const workerTasksUsage = this.workerNodes[workerNodeKey].tasksUsage
--workerTasksUsage.running
++workerTasksUsage.run
if (message.error != null) {
workerTasksUsage.avgRunTime =
workerTasksUsage.runTime / workerTasksUsage.run
}
- if (this.workerChoiceStrategyContext.getRequiredStatistics().medRunTime) {
- workerTasksUsage.runTimeHistory.push(message.runTime ?? 0)
+ if (
+ this.workerChoiceStrategyContext.getRequiredStatistics().medRunTime &&
+ message.runTime != null
+ ) {
+ workerTasksUsage.runTimeHistory.push(message.runTime)
workerTasksUsage.medRunTime = median(workerTasksUsage.runTimeHistory)
}
}
- this.workerChoiceStrategyContext.update()
}
/**
* Chooses a worker node for the next task.
*
- * The default uses a round robin algorithm to distribute the load.
+ * The default worker choice strategy uses a round robin algorithm to distribute the load.
*
- * @returns [worker node key, worker node].
+ * @returns The worker node key
*/
- protected chooseWorkerNode (): [number, WorkerNode<Worker, Data>] {
+ protected chooseWorkerNode (): number {
let workerNodeKey: number
if (this.type === PoolType.DYNAMIC && !this.full && this.internalBusy()) {
const workerCreated = this.createAndSetupWorker()
this.registerWorkerMessageListener(workerCreated, message => {
+ const currentWorkerNodeKey = this.getWorkerNodeKey(workerCreated)
if (
isKillBehavior(KillBehaviors.HARD, message.kill) ||
(message.kill != null &&
- this.getWorkerTasksUsage(workerCreated)?.running === 0)
+ this.workerNodes[currentWorkerNodeKey].tasksUsage.running === 0)
) {
// Kill message received from the worker: no new tasks are submitted to that worker for a while ( > maxInactiveTime)
- this.flushTasksQueueByWorker(workerCreated)
+ this.flushTasksQueue(currentWorkerNodeKey)
void (this.destroyWorker(workerCreated) as Promise<void>)
}
})
} else {
workerNodeKey = this.workerChoiceStrategyContext.execute()
}
- return [workerNodeKey, this.workerNodes[workerNodeKey]]
+ return workerNodeKey
}
/**
workerNode.tasksUsage = tasksUsage
}
- /**
- * Gets the given worker its tasks usage in the pool.
- *
- * @param worker - The worker.
- * @throws Error if the worker is not found in the pool worker nodes.
- * @returns The worker tasks usage.
- */
- private getWorkerTasksUsage (worker: Worker): TasksUsage {
- const workerNodeKey = this.getWorkerNodeKey(worker)
- if (workerNodeKey !== -1) {
- return this.workerNodes[workerNodeKey].tasksUsage
- }
- throw new Error('Worker could not be found in the pool worker nodes')
- }
-
/**
* Pushes the given worker in the pool worker nodes.
*
}
}
- private flushTasksQueueByWorker (worker: Worker): void {
- const workerNodeKey = this.getWorkerNodeKey(worker)
- this.flushTasksQueue(workerNodeKey)
- }
-
private flushTasksQueues (): void {
for (const [workerNodeKey] of this.workerNodes.entries()) {
this.flushTasksQueue(workerNodeKey)