import { AsyncResource } from 'node:async_hooks'
import type { Worker } from 'node:cluster'
import type { MessagePort } from 'node:worker_threads'
-import { type EventLoopUtilization, performance } from 'node:perf_hooks'
-import type { MessageValue, WorkerStatistics } from '../utility-types'
+import { performance } from 'node:perf_hooks'
+import type {
+ MessageValue,
+ TaskPerformance,
+ WorkerStatistics
+} from '../utility-types'
import { EMPTY_FUNCTION, isPlainObject } from '../utils'
import {
type KillBehavior,
const DEFAULT_MAX_INACTIVE_TIME = 60000
const DEFAULT_KILL_BEHAVIOR: KillBehavior = KillBehaviors.SOFT
-/**
- * Task performance.
- */
-export interface TaskPerformance {
- timestamp: number
- waitTime?: number
- runTime?: number
- elu?: EventLoopUtilization
-}
-
/**
* Base class that implements some shared logic for all poolifier workers.
*
*/
protected lastTaskTimestamp!: number
/**
- * Performance statistics computation.
+ * Performance statistics computation requirements.
*/
protected statistics!: WorkerStatistics
/**
*
* @param message - Message received.
*/
- protected messageListener (message: MessageValue<Data, MainWorker>): void {
+ protected messageListener (
+ message: MessageValue<Data, Data, MainWorker>
+ ): void {
if (message.id != null && message.data != null) {
// Task message received
const fn = this.getTaskFunction(message.name)
} else if (message.parent != null) {
// Main worker reference message received
this.mainWorker = message.parent
+ } else if (message.statistics != null) {
+ // Statistics message received
+ this.statistics = message.statistics
} else if (message.kill != null) {
// Kill message received
this.aliveInterval != null && clearInterval(this.aliveInterval)
this.emitDestroy()
- } else if (message.statistics != null) {
- // Statistics message received
- this.statistics = message.statistics
}
}
*
* @param message - The response message.
*/
- protected abstract sendToMainWorker (message: MessageValue<Response>): void
+ protected abstract sendToMainWorker (
+ message: MessageValue<Response, Data>
+ ): void
/**
* Checks if the worker should be terminated, because its living too long.
message: MessageValue<Data>
): void {
try {
- const taskPerformance = this.beginTaskPerformance(message)
+ let taskPerformance = this.beginTaskPerformance()
const res = fn(message.data)
- const { runTime, waitTime, elu } =
- this.endTaskPerformance(taskPerformance)
+ taskPerformance = this.endTaskPerformance(taskPerformance)
this.sendToMainWorker({
data: res,
- runTime,
- waitTime,
- elu,
+ taskPerformance,
id: message.id
})
} catch (e) {
const err = this.handleError(e as Error)
this.sendToMainWorker({
- error: err,
- errorData: message.data,
+ taskError: {
+ message: err,
+ data: message.data
+ },
id: message.id
})
} finally {
fn: WorkerAsyncFunction<Data, Response>,
message: MessageValue<Data>
): void {
- const taskPerformance = this.beginTaskPerformance(message)
+ let taskPerformance = this.beginTaskPerformance()
fn(message.data)
.then(res => {
- const { runTime, waitTime, elu } =
- this.endTaskPerformance(taskPerformance)
+ taskPerformance = this.endTaskPerformance(taskPerformance)
this.sendToMainWorker({
data: res,
- runTime,
- waitTime,
- elu,
+ taskPerformance,
id: message.id
})
return null
.catch(e => {
const err = this.handleError(e as Error)
this.sendToMainWorker({
- error: err,
- errorData: message.data,
+ taskError: {
+ message: err,
+ data: message.data
+ },
id: message.id
})
})
return fn
}
- private beginTaskPerformance (message: MessageValue<Data>): TaskPerformance {
- const timestamp = performance.now()
+ private beginTaskPerformance (): TaskPerformance {
+ this.checkStatistics()
return {
- timestamp,
- ...(this.statistics.waitTime && {
- waitTime: timestamp - (message.timestamp ?? timestamp)
- }),
+ timestamp: performance.now(),
...(this.statistics.elu && { elu: performance.eventLoopUtilization() })
}
}
private endTaskPerformance (
taskPerformance: TaskPerformance
): TaskPerformance {
+ this.checkStatistics()
return {
...taskPerformance,
...(this.statistics.runTime && {
})
}
}
+
+ private checkStatistics (): void {
+ if (this.statistics == null) {
+ throw new Error('Performance statistics computation requirements not set')
+ }
+ }
}