repositories
/
poolifier.git
/ blobdiff
commit
grep
author
committer
pickaxe
?
search:
re
summary
|
shortlog
|
log
|
commit
|
commitdiff
|
tree
raw
|
inline
| side by side
Merge dependabot/npm_and_yarn/examples/typescript/http-server-pool/fastify-hybrid...
[poolifier.git]
/
src
/
pools
/
worker.ts
diff --git
a/src/pools/worker.ts
b/src/pools/worker.ts
index e4a31e393fad35e06368e3d33cc5820ca1dbd099..5439606d420f0abd165b84c49157808de230aee6 100644
(file)
--- a/
src/pools/worker.ts
+++ b/
src/pools/worker.ts
@@
-1,14
+1,19
@@
import type { MessageChannel } from 'node:worker_threads'
import type { MessageChannel } from 'node:worker_threads'
+import type { EventEmitter } from 'node:events'
import type { CircularArray } from '../circular-array'
import type { Task } from '../utility-types'
/**
* Callback invoked when the worker has started successfully.
import type { CircularArray } from '../circular-array'
import type { Task } from '../utility-types'
/**
* 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.
*/
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,
*/
export type MessageHandler<Worker extends IWorker> = (
this: Worker,
@@
-17,6
+22,8
@@
export type MessageHandler<Worker extends IWorker> = (
/**
* Callback invoked if the worker raised an error.
/**
* Callback invoked if the worker raised an error.
+ *
+ * @typeParam Worker - Type of worker.
*/
export type ErrorHandler<Worker extends IWorker> = (
this: Worker,
*/
export type ErrorHandler<Worker extends IWorker> = (
this: Worker,
@@
-25,6
+32,8
@@
export type ErrorHandler<Worker extends IWorker> = (
/**
* Callback invoked when the worker exits successfully.
/**
* Callback invoked when the worker exits successfully.
+ *
+ * @typeParam Worker - Type of worker.
*/
export type ExitHandler<Worker extends IWorker> = (
this: Worker,
*/
export type ExitHandler<Worker extends IWorker> = (
this: Worker,
@@
-96,6
+105,14
@@
export interface TaskStatistics {
* Maximum number of queued tasks.
*/
readonly maxQueued?: number
* Maximum number of queued tasks.
*/
readonly maxQueued?: number
+ /**
+ * Number of sequentially stolen tasks.
+ */
+ sequentiallyStolen: number
+ /**
+ * Number of stolen tasks.
+ */
+ stolen: number
/**
* Number of failed tasks.
*/
/**
* Number of failed tasks.
*/
@@
-128,7
+145,7
@@
export interface WorkerInfo {
/**
* Worker type.
*/
/**
* Worker type.
*/
- type: WorkerType
+
readonly
type: WorkerType
/**
* Dynamic flag.
*/
/**
* Dynamic flag.
*/
@@
-140,7
+157,7
@@
export interface WorkerInfo {
/**
* Task function names.
*/
/**
* Task function names.
*/
- taskFunctions?: string[]
+ taskFunction
Name
s?: string[]
}
/**
}
/**
@@
-167,6
+184,15
@@
export interface WorkerUsage {
readonly elu: EventLoopUtilizationMeasurementStatistics
}
readonly elu: EventLoopUtilizationMeasurementStatistics
}
+/**
+ * Worker choice strategy data.
+ *
+ * @internal
+ */
+export interface StrategyData {
+ virtualTaskEndTimestamp?: number
+}
+
/**
* Worker interface.
*/
/**
* Worker interface.
*/
@@
-189,12
+215,22
@@
export interface IWorker {
/**
* Registers a listener to the exit event that will only be performed once.
*
/**
* Registers a listener to the exit event that will only be performed once.
*
- * @param event -
`'exit'`
.
+ * @param event -
The `'exit'` event
.
* @param handler - The exit handler.
*/
readonly once: (event: 'exit', handler: ExitHandler<this>) => void
}
* @param handler - The exit handler.
*/
readonly once: (event: 'exit', handler: ExitHandler<this>) => void
}
+/**
+ * Worker node event detail.
+ *
+ * @internal
+ */
+export interface WorkerNodeEventDetail {
+ workerId: number
+ workerNodeKey?: number
+}
+
/**
* Worker node interface.
*
/**
* Worker node interface.
*
@@
-202,7
+238,8
@@
export interface IWorker {
* @typeParam Data - Type of data sent to the worker. This can only be structured-cloneable data.
* @internal
*/
* @typeParam Data - Type of data sent to the worker. This can only be structured-cloneable data.
* @internal
*/
-export interface IWorkerNode<Worker extends IWorker, Data = unknown> {
+export interface IWorkerNode<Worker extends IWorker, Data = unknown>
+ extends EventEmitter {
/**
* Worker.
*/
/**
* Worker.
*/
@@
-211,14
+248,24
@@
export interface IWorkerNode<Worker extends IWorker, Data = unknown> {
* Worker info.
*/
readonly info: WorkerInfo
* 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_threads only).
*/
readonly messageChannel?: MessageChannel
/**
/**
* Message channel (worker_threads only).
*/
readonly messageChannel?: MessageChannel
/**
- * Worker usage statistics.
+ * Tasks queue back pressure size.
+ * This is the number of tasks that can be enqueued before the worker node has back pressure.
*/
*/
- usage: WorkerUsage
+ tasksQueueBackPressureSize: number
/**
* Tasks queue size.
*
/**
* Tasks queue size.
*
@@
-232,12
+279,25
@@
export interface IWorkerNode<Worker extends IWorker, Data = unknown> {
* @returns The tasks queue size.
*/
readonly enqueueTask: (task: Task<Data>) => number
* @returns The tasks queue size.
*/
readonly enqueueTask: (task: Task<Data>) => number
+ /**
+ * Prepends a task to the tasks queue.
+ *
+ * @param task - The task to prepend.
+ * @returns The tasks queue size.
+ */
+ readonly unshiftTask: (task: Task<Data>) => number
/**
* Dequeue task.
*
* @returns The dequeued task.
*/
readonly dequeueTask: () => Task<Data> | undefined
/**
* Dequeue task.
*
* @returns The dequeued task.
*/
readonly dequeueTask: () => Task<Data> | undefined
+ /**
+ * Pops a task from the tasks queue.
+ *
+ * @returns The popped task.
+ */
+ readonly popTask: () => Task<Data> | undefined
/**
* Clears tasks queue.
*/
/**
* Clears tasks queue.
*/
@@
-263,4
+323,11
@@
export interface IWorkerNode<Worker extends IWorker, Data = unknown> {
* @returns The task function worker usage statistics if the task function worker usage statistics are initialized, `undefined` otherwise.
*/
readonly getTaskFunctionWorkerUsage: (name: string) => WorkerUsage | undefined
* @returns The task function worker usage statistics if the task function worker usage statistics are initialized, `undefined` otherwise.
*/
readonly getTaskFunctionWorkerUsage: (name: string) => WorkerUsage | undefined
+ /**
+ * 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 deleteTaskFunctionWorkerUsage: (name: string) => boolean
}
}