median
} from '../utils'
import { KillBehaviors, isKillBehavior } from '../worker/worker-options'
+import { CircularArray } from '../circular-array'
+import { Queue } from '../queue'
import {
- PoolEvents,
type IPool,
+ PoolEmitter,
+ PoolEvents,
type PoolOptions,
- type TasksQueueOptions,
- PoolType
+ PoolType,
+ type TasksQueueOptions
} from './pool'
-import { PoolEmitter } from './pool'
import type { IWorker, Task, TasksUsage, WorkerNode } from './worker'
import {
WorkerChoiceStrategies,
type WorkerChoiceStrategyOptions
} from './selection-strategies/selection-strategies-types'
import { WorkerChoiceStrategyContext } from './selection-strategies/worker-choice-strategy-context'
-import { CircularArray } from '../circular-array'
/**
* Base class that implements some shared logic for all poolifier pools.
* Constructs a new poolifier pool.
*
* @param numberOfWorkers - Number of workers that this pool should manage.
- * @param filePath - Path to the worker-file.
+ * @param filePath - Path to the worker file.
* @param opts - Options for the pool.
*/
public constructor (
this.checkFilePath(this.filePath)
this.checkPoolOptions(this.opts)
- this.chooseWorkerNode.bind(this)
- this.executeTask.bind(this)
- this.enqueueTask.bind(this)
- this.checkAndEmitEvents.bind(this)
+ this.chooseWorkerNode = this.chooseWorkerNode.bind(this)
+ this.executeTask = this.executeTask.bind(this)
+ this.enqueueTask = this.enqueueTask.bind(this)
+ this.checkAndEmitEvents = this.checkAndEmitEvents.bind(this)
this.setupHook()
return 0
}
return this.workerNodes.reduce(
- (accumulator, workerNode) => accumulator + workerNode.tasksQueue.length,
+ (accumulator, workerNode) => accumulator + workerNode.tasksQueue.size,
0
)
}
/** @inheritDoc */
public setWorkerChoiceStrategy (
- workerChoiceStrategy: WorkerChoiceStrategy
+ workerChoiceStrategy: WorkerChoiceStrategy,
+ workerChoiceStrategyOptions?: WorkerChoiceStrategyOptions
): void {
this.checkValidWorkerChoiceStrategy(workerChoiceStrategy)
this.opts.workerChoiceStrategy = workerChoiceStrategy
this.workerChoiceStrategyContext.setWorkerChoiceStrategy(
this.opts.workerChoiceStrategy
)
+ if (workerChoiceStrategyOptions != null) {
+ this.setWorkerChoiceStrategyOptions(workerChoiceStrategyOptions)
+ }
}
/** @inheritDoc */
}
/** @inheritDoc */
- public enableTasksQueue (enable: boolean, opts?: TasksQueueOptions): void {
+ public enableTasksQueue (
+ enable: boolean,
+ tasksQueueOptions?: TasksQueueOptions
+ ): void {
if (this.opts.enableTasksQueue === true && !enable) {
- for (const [workerNodeKey] of this.workerNodes.entries()) {
- this.flushTasksQueue(workerNodeKey)
- }
+ this.flushTasksQueues()
}
this.opts.enableTasksQueue = enable
- this.setTasksQueueOptions(opts as TasksQueueOptions)
+ this.setTasksQueueOptions(tasksQueueOptions as TasksQueueOptions)
}
/** @inheritDoc */
- public setTasksQueueOptions (opts: TasksQueueOptions): void {
+ public setTasksQueueOptions (tasksQueueOptions: TasksQueueOptions): void {
if (this.opts.enableTasksQueue === true) {
- this.checkValidTasksQueueOptions(opts)
- this.opts.tasksQueueOptions = this.buildTasksQueueOptions(opts)
+ this.checkValidTasksQueueOptions(tasksQueueOptions)
+ this.opts.tasksQueueOptions =
+ this.buildTasksQueueOptions(tasksQueueOptions)
} else {
delete this.opts.tasksQueueOptions
}
protected abstract get busy (): boolean
protected internalBusy (): boolean {
- return this.findFreeWorkerNodeKey() === -1
- }
-
- /** @inheritDoc */
- public findFreeWorkerNodeKey (): number {
- return this.workerNodes.findIndex(workerNode => {
- return workerNode.tasksUsage?.running === 0
- })
+ return (
+ this.workerNodes.findIndex(workerNode => {
+ return workerNode.tasksUsage?.running === 0
+ }) === -1
+ )
}
/** @inheritDoc */
- public async execute (data: Data): Promise<Response> {
+ public async execute (data?: Data): Promise<Response> {
const [workerNodeKey, workerNode] = this.chooseWorkerNode()
const submittedTask: Task<Data> = {
// eslint-disable-next-line @typescript-eslint/consistent-type-assertions
worker: Worker,
message: MessageValue<Response>
): void {
- const workerTasksUsage = this.getWorkerTasksUsage(worker) as TasksUsage
+ const workerTasksUsage = this.getWorkerTasksUsage(worker)
--workerTasksUsage.running
++workerTasksUsage.run
if (message.error != null) {
) {
// Kill message received from the worker: no new tasks are submitted to that worker for a while ( > maxInactiveTime)
this.flushTasksQueueByWorker(workerCreated)
- void this.destroyWorker(workerCreated)
+ void (this.destroyWorker(workerCreated) as Promise<void>)
}
})
workerNodeKey = this.getWorkerNodeKey(workerCreated)
* Gets the given worker its tasks usage in the pool.
*
* @param worker - The worker.
- * @throws {@link Error} if the worker is not found in the pool worker nodes.
+ * @throws Error if the worker is not found in the pool worker nodes.
* @returns The worker tasks usage.
*/
- private getWorkerTasksUsage (worker: Worker): TasksUsage | undefined {
+ private getWorkerTasksUsage (worker: Worker): TasksUsage {
const workerNodeKey = this.getWorkerNodeKey(worker)
if (workerNodeKey !== -1) {
return this.workerNodes[workerNodeKey].tasksUsage
medRunTime: 0,
error: 0
},
- tasksQueue: []
+ tasksQueue: new Queue<Task<Data>>()
})
}
workerNodeKey: number,
worker: Worker,
tasksUsage: TasksUsage,
- tasksQueue: Array<Task<Data>>
+ tasksQueue: Queue<Task<Data>>
): void {
this.workerNodes[workerNodeKey] = {
worker,
}
private enqueueTask (workerNodeKey: number, task: Task<Data>): number {
- return this.workerNodes[workerNodeKey].tasksQueue.push(task)
+ return this.workerNodes[workerNodeKey].tasksQueue.enqueue(task)
}
private dequeueTask (workerNodeKey: number): Task<Data> | undefined {
- return this.workerNodes[workerNodeKey].tasksQueue.shift()
+ return this.workerNodes[workerNodeKey].tasksQueue.dequeue()
}
private tasksQueueSize (workerNodeKey: number): number {
- return this.workerNodes[workerNodeKey].tasksQueue.length
+ return this.workerNodes[workerNodeKey].tasksQueue.size
}
private flushTasksQueue (workerNodeKey: number): void {
if (this.tasksQueueSize(workerNodeKey) > 0) {
- for (const task of this.workerNodes[workerNodeKey].tasksQueue) {
- this.executeTask(workerNodeKey, task)
+ for (let i = 0; i < this.tasksQueueSize(workerNodeKey); i++) {
+ this.executeTask(
+ workerNodeKey,
+ this.dequeueTask(workerNodeKey) as Task<Data>
+ )
}
}
}
const workerNodeKey = this.getWorkerNodeKey(worker)
this.flushTasksQueue(workerNodeKey)
}
+
+ private flushTasksQueues (): void {
+ for (const [workerNodeKey] of this.workerNodes.entries()) {
+ this.flushTasksQueue(workerNodeKey)
+ }
+ }
}