import {
DEFAULT_WORKER_CHOICE_STRATEGY_OPTIONS,
EMPTY_FUNCTION,
+ isPlainObject,
median
} from '../utils'
import { KillBehaviors, isKillBehavior } from '../worker/worker-options'
)
} else if (!Number.isSafeInteger(numberOfWorkers)) {
throw new TypeError(
- 'Cannot instantiate a pool with a non integer number of workers'
+ 'Cannot instantiate a pool with a non safe integer number of workers'
)
} else if (numberOfWorkers < 0) {
throw new RangeError(
}
private checkPoolOptions (opts: PoolOptions<Worker>): void {
- this.opts.workerChoiceStrategy =
- opts.workerChoiceStrategy ?? WorkerChoiceStrategies.ROUND_ROBIN
- this.checkValidWorkerChoiceStrategy(this.opts.workerChoiceStrategy)
- this.opts.workerChoiceStrategyOptions =
- opts.workerChoiceStrategyOptions ?? DEFAULT_WORKER_CHOICE_STRATEGY_OPTIONS
- this.opts.enableEvents = opts.enableEvents ?? true
- this.opts.enableTasksQueue = opts.enableTasksQueue ?? false
- if (this.opts.enableTasksQueue) {
- this.checkValidTasksQueueOptions(
- opts.tasksQueueOptions as TasksQueueOptions
- )
- this.opts.tasksQueueOptions = this.buildTasksQueueOptions(
- opts.tasksQueueOptions as TasksQueueOptions
+ if (isPlainObject(opts)) {
+ this.opts.workerChoiceStrategy =
+ opts.workerChoiceStrategy ?? WorkerChoiceStrategies.ROUND_ROBIN
+ this.checkValidWorkerChoiceStrategy(this.opts.workerChoiceStrategy)
+ 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) {
+ this.checkValidTasksQueueOptions(
+ opts.tasksQueueOptions as TasksQueueOptions
+ )
+ this.opts.tasksQueueOptions = this.buildTasksQueueOptions(
+ opts.tasksQueueOptions as TasksQueueOptions
+ )
+ }
+ } else {
+ throw new TypeError('Invalid pool options: must be a plain object')
}
}
}
}
+ private checkValidWorkerChoiceStrategyOptions (
+ workerChoiceStrategyOptions: WorkerChoiceStrategyOptions
+ ): void {
+ if (!isPlainObject(workerChoiceStrategyOptions)) {
+ throw new TypeError(
+ '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 (
tasksQueueOptions: TasksQueueOptions
): void {
+ if (tasksQueueOptions != null && !isPlainObject(tasksQueueOptions)) {
+ throw new TypeError('Invalid tasks queue options: must be a plain object')
+ }
if ((tasksQueueOptions?.concurrency as number) <= 0) {
throw new Error(
`Invalid worker tasks concurrency '${
/** @inheritDoc */
public abstract get type (): PoolType
+ /** @inheritDoc */
+ public abstract get size (): number
+
/**
* Number of tasks running in the pool.
*/
runTimeHistory: new CircularArray(),
avgRunTime: 0,
medRunTime: 0,
+ waitTime: 0,
+ waitTimeHistory: new CircularArray(),
+ avgWaitTime: 0,
+ medWaitTime: 0,
error: 0
})
}
public setWorkerChoiceStrategyOptions (
workerChoiceStrategyOptions: WorkerChoiceStrategyOptions
): void {
+ this.checkValidWorkerChoiceStrategyOptions(workerChoiceStrategyOptions)
this.opts.workerChoiceStrategyOptions = workerChoiceStrategyOptions
this.workerChoiceStrategyContext.setOptions(
this.opts.workerChoiceStrategyOptions
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 submissionTimestamp = performance.now()
+ const workerNodeKey = this.chooseWorkerNode()
const submittedTask: Task<Data> = {
name,
// eslint-disable-next-line @typescript-eslint/consistent-type-assertions
data: data ?? ({} as Data),
+ submissionTimestamp,
id: crypto.randomUUID()
}
const res = new Promise<Response>((resolve, reject) => {
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 workerTasksUsage =
+ this.workerNodes[this.getWorkerNodeKey(worker)].tasksUsage
--workerTasksUsage.running
++workerTasksUsage.run
if (message.error != null) {
++workerTasksUsage.error
}
+ this.updateRunTimeTasksUsage(workerTasksUsage, message)
+ this.updateWaitTasksUsage(workerTasksUsage, message)
+ }
+
+ private updateRunTimeTasksUsage (
+ workerTasksUsage: TasksUsage,
+ message: MessageValue<Response>
+ ): void {
if (this.workerChoiceStrategyContext.getRequiredStatistics().runTime) {
workerTasksUsage.runTime += message.runTime ?? 0
if (
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)
}
}
}
+ private updateWaitTasksUsage (
+ workerTasksUsage: TasksUsage,
+ message: MessageValue<Response>
+ ): void {
+ if (this.workerChoiceStrategyContext.getRequiredStatistics().waitTime) {
+ workerTasksUsage.waitTime += message.waitTime ?? 0
+ if (
+ this.workerChoiceStrategyContext.getRequiredStatistics().avgWaitTime &&
+ workerTasksUsage.run !== 0
+ ) {
+ workerTasksUsage.avgWaitTime =
+ workerTasksUsage.waitTime / workerTasksUsage.run
+ }
+ if (
+ this.workerChoiceStrategyContext.getRequiredStatistics().medWaitTime &&
+ message.waitTime != null
+ ) {
+ workerTasksUsage.waitTimeHistory.push(message.waitTime)
+ workerTasksUsage.medWaitTime = median(workerTasksUsage.waitTimeHistory)
+ }
+ }
+ }
+
/**
* 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.
*
runTimeHistory: new CircularArray(),
avgRunTime: 0,
medRunTime: 0,
+ waitTime: 0,
+ waitTimeHistory: new CircularArray(),
+ avgWaitTime: 0,
+ medWaitTime: 0,
error: 0
},
tasksQueue: new Queue<Task<Data>>()
}
}
- 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)