-import type { EventLoopUtilization } from 'node:perf_hooks'
-import type { CircularArray } from '../circular-array'
-import type { Queue } from '../queue'
+import type { EventEmitter } from 'node:events'
+import type { MessageChannel, WorkerOptions } from 'node:worker_threads'
+
+import type { CircularArray } from '../circular-array.js'
+import type { Task, TaskFunctionProperties } from '../utility-types.js'
+
+/**
+ * Callback invoked when the worker has started successfully.
+ *
+ * @typeParam Worker - Type of worker.
+ */
+export type OnlineHandler<Worker extends IWorker> = (this: Worker) => void
/**
* Callback invoked if the worker has received a message.
+ *
+ * @typeParam Worker - Type of worker.
*/
export type MessageHandler<Worker extends IWorker> = (
this: Worker,
- m: unknown
+ message: unknown
) => void
/**
* Callback invoked if the worker raised an error.
+ *
+ * @typeParam Worker - Type of worker.
*/
export type ErrorHandler<Worker extends IWorker> = (
this: Worker,
- e: Error
+ error: Error
) => void
-/**
- * Callback invoked when the worker has started successfully.
- */
-export type OnlineHandler<Worker extends IWorker> = (this: Worker) => void
-
/**
* Callback invoked when the worker exits successfully.
+ *
+ * @typeParam Worker - Type of worker.
*/
export type ExitHandler<Worker extends IWorker> = (
this: Worker,
- code: number
+ exitCode: number
) => void
/**
- * Message object that is passed as a task between main worker and worker.
+ * Worker event handler.
*
- * @typeParam Data - Type of data sent to the worker. This can only be serializable data.
- * @internal
+ * @typeParam Worker - Type of worker.
*/
-export interface Task<Data = unknown> {
- /**
- * Task name.
- */
- readonly name?: string
- /**
- * Task input data that will be passed to the worker.
- */
- readonly data?: Data
- /**
- * Timestamp.
- */
- readonly timestamp?: number
- /**
- * Message UUID.
- */
- readonly id?: string
-}
+export type EventHandler<Worker extends IWorker> =
+ | OnlineHandler<Worker>
+ | MessageHandler<Worker>
+ | ErrorHandler<Worker>
+ | ExitHandler<Worker>
/**
* Measurement statistics.
*/
export interface MeasurementStatistics {
/**
- * Measurement aggregation.
+ * Measurement aggregate.
*/
- aggregation: number
+ aggregate?: number
+ /**
+ * Measurement minimum.
+ */
+ minimum?: number
+ /**
+ * Measurement maximum.
+ */
+ maximum?: number
/**
* Measurement average.
*/
- average: number
+ average?: number
/**
* Measurement median.
*/
- median: number
+ median?: number
/**
* Measurement history.
*/
- history: CircularArray<number>
+ readonly history: CircularArray<number>
+}
+
+/**
+ * Event loop utilization measurement statistics.
+ *
+ * @internal
+ */
+export interface EventLoopUtilizationMeasurementStatistics {
+ readonly idle: MeasurementStatistics
+ readonly active: MeasurementStatistics
+ utilization?: number
}
/**
*/
export interface TaskStatistics {
/**
- * Number of tasks executed.
+ * Number of executed tasks.
*/
executed: number
/**
- * Number of tasks executing.
+ * Number of executing tasks.
*/
executing: number
/**
- * Number of tasks queued.
+ * Number of queued tasks.
+ */
+ readonly queued: number
+ /**
+ * Maximum number of queued tasks.
+ */
+ readonly maxQueued?: number
+ /**
+ * Number of sequentially stolen tasks.
+ */
+ sequentiallyStolen: number
+ /**
+ * Number of stolen tasks.
*/
- queued: number
+ stolen: number
/**
- * Number of tasks failed.
+ * Number of failed tasks.
*/
failed: number
}
+/**
+ * Enumeration of worker types.
+ */
+export const WorkerTypes: Readonly<{ thread: 'thread', cluster: 'cluster' }> =
+ Object.freeze({
+ thread: 'thread',
+ cluster: 'cluster'
+ } as const)
+
+/**
+ * Worker type.
+ */
+export type WorkerType = keyof typeof WorkerTypes
+
+/**
+ * Worker information.
+ *
+ * @internal
+ */
+export interface WorkerInfo {
+ /**
+ * Worker id.
+ */
+ readonly id: number | undefined
+ /**
+ * Worker type.
+ */
+ readonly type: WorkerType
+ /**
+ * Dynamic flag.
+ */
+ dynamic: boolean
+ /**
+ * Ready flag.
+ */
+ ready: boolean
+ /**
+ * Stealing flag.
+ * This flag is set to `true` when worker node is stealing tasks from another worker node.
+ */
+ stealing: boolean
+ /**
+ * Back pressure flag.
+ * This flag is set to `true` when worker node tasks queue has back pressure.
+ */
+ backPressure: boolean
+ /**
+ * Task functions properties.
+ */
+ taskFunctionsProperties?: TaskFunctionProperties[]
+}
+
/**
* Worker usage statistics.
*
/**
* Tasks statistics.
*/
- tasks: TaskStatistics
+ readonly tasks: TaskStatistics
/**
* Tasks runtime statistics.
*/
- runTime: MeasurementStatistics
+ readonly runTime: MeasurementStatistics
/**
* Tasks wait time statistics.
*/
- waitTime: MeasurementStatistics
+ readonly waitTime: MeasurementStatistics
/**
- * Event loop utilization.
+ * Tasks event loop utilization statistics.
*/
- elu: EventLoopUtilization | undefined
+ readonly elu: EventLoopUtilizationMeasurementStatistics
+}
+
+/**
+ * Worker choice strategy data.
+ *
+ * @internal
+ */
+export interface StrategyData {
+ virtualTaskEndTimestamp?: number
}
/**
* Worker interface.
*/
-export interface IWorker {
+export interface IWorker extends EventEmitter {
+ /**
+ * Cluster worker id.
+ */
+ readonly id?: number
+ /**
+ * Worker thread worker id.
+ */
+ readonly threadId?: number
/**
- * Register an event listener.
+ * Registers an event handler.
*
* @param event - The event.
* @param handler - The event handler.
*/
- on: ((event: 'message', handler: MessageHandler<this>) => void) &
- ((event: 'error', handler: ErrorHandler<this>) => void) &
- ((event: 'online', handler: OnlineHandler<this>) => void) &
- ((event: 'exit', handler: ExitHandler<this>) => void)
+ readonly on: (event: string, handler: EventHandler<this>) => this
/**
- * Register a listener to the exit event that will only be performed once.
+ * Registers once an event handler.
*
- * @param event - `'exit'`.
- * @param handler - The exit handler.
+ * @param event - The event.
+ * @param handler - The event handler.
+ */
+ readonly once: (event: string, handler: EventHandler<this>) => this
+ /**
+ * Calling `unref()` on a worker allows the thread to exit if this is the only
+ * active handle in the event system. If the worker is already `unref()`ed calling`unref()` again has no effect.
+ * @since v10.5.0
+ */
+ readonly unref?: () => void
+ /**
+ * Stop all JavaScript execution in the worker thread as soon as possible.
+ * Returns a Promise for the exit code that is fulfilled when the `'exit' event` is emitted.
+ */
+ readonly terminate?: () => Promise<number>
+ /**
+ * Cluster worker disconnect.
+ */
+ readonly disconnect?: () => void
+ /**
+ * Cluster worker kill.
*/
- once: (event: 'exit', handler: ExitHandler<this>) => void
+ readonly kill?: (signal?: string) => void
+}
+
+/**
+ * Worker node options.
+ *
+ * @internal
+ */
+export interface WorkerNodeOptions {
+ workerOptions?: WorkerOptions
+ env?: Record<string, unknown>
+ tasksQueueBackPressureSize: number | undefined
+ tasksQueueBucketSize: number | undefined
}
/**
* Worker node interface.
*
* @typeParam Worker - Type of worker.
- * @typeParam Data - Type of data sent to the worker. This can only be serializable data.
+ * @typeParam Data - Type of data sent to the worker. This can only be structured-cloneable data.
* @internal
*/
-export interface WorkerNode<Worker extends IWorker, Data = unknown> {
+export interface IWorkerNode<Worker extends IWorker, Data = unknown>
+ extends EventEmitter {
/**
- * Worker node worker.
+ * Worker.
*/
readonly worker: Worker
/**
- * Worker node worker usage statistics.
+ * Worker info.
+ */
+ readonly info: WorkerInfo
+ /**
+ * Worker usage statistics.
+ */
+ readonly usage: WorkerUsage
+ /**
+ * Worker choice strategy data.
+ * This is used to store data that are specific to the worker choice strategy.
+ */
+ strategyData?: StrategyData
+ /**
+ * Message channel (worker thread only).
+ */
+ readonly messageChannel?: MessageChannel
+ /**
+ * Tasks queue back pressure size.
+ * This is the number of tasks that can be enqueued before the worker node has back pressure.
+ */
+ tasksQueueBackPressureSize: number
+ /**
+ * Tasks queue size.
+ *
+ * @returns The tasks queue size.
+ */
+ readonly tasksQueueSize: () => number
+ /**
+ * Enqueue task.
+ *
+ * @param task - The task to queue.
+ * @returns The tasks queue size.
+ */
+ readonly enqueueTask: (task: Task<Data>) => number
+ /**
+ * Dequeue task.
+ *
+ * @param bucket - The prioritized bucket to dequeue from. @defaultValue 0
+ * @returns The dequeued task.
+ */
+ readonly dequeueTask: (bucket?: number) => Task<Data> | undefined
+ /**
+ * Dequeue last prioritized task.
+ *
+ * @returns The dequeued task.
+ */
+ readonly dequeueLastPrioritizedTask: () => Task<Data> | undefined
+ /**
+ * Clears tasks queue.
+ */
+ readonly clearTasksQueue: () => void
+ /**
+ * Whether the worker node has back pressure (i.e. its tasks queue is full).
+ *
+ * @returns `true` if the worker node has back pressure, `false` otherwise.
+ */
+ readonly hasBackPressure: () => boolean
+ /**
+ * Terminates the worker node.
+ */
+ readonly terminate: () => Promise<void>
+ /**
+ * Registers a worker event handler.
+ *
+ * @param event - The event.
+ * @param handler - The event handler.
+ */
+ readonly registerWorkerEventHandler: (
+ event: string,
+ handler: EventHandler<Worker>
+ ) => void
+ /**
+ * Registers once a worker event handler.
+ *
+ * @param event - The event.
+ * @param handler - The event handler.
+ */
+ readonly registerOnceWorkerEventHandler: (
+ event: string,
+ handler: EventHandler<Worker>
+ ) => void
+ /**
+ * Gets task function worker usage statistics.
+ *
+ * @param name - The task function name.
+ * @returns The task function worker usage statistics if the task function worker usage statistics are initialized, `undefined` otherwise.
*/
- workerUsage: WorkerUsage
+ readonly getTaskFunctionWorkerUsage: (name: string) => WorkerUsage | undefined
/**
- * Worker node tasks queue.
+ * Deletes task function worker usage statistics.
+ *
+ * @param name - The task function name.
+ * @returns `true` if the task function worker usage statistics were deleted, `false` otherwise.
*/
- readonly tasksQueue: Queue<Task<Data>>
+ readonly deleteTaskFunctionWorkerUsage: (name: string) => boolean
+}
+
+/**
+ * Worker node event detail.
+ *
+ * @internal
+ */
+export interface WorkerNodeEventDetail {
+ workerId?: number
+ workerNodeKey?: number
}