Commit | Line | Data |
---|---|---|
962a8159 | 1 | import type EventEmitterAsyncResource from 'node:events'; |
01f4001e | 2 | import type { Worker } from 'node:worker_threads'; |
0e4fa348 | 3 | |
b779c0f8 | 4 | import { DynamicThreadPool, type ErrorHandler, type ExitHandler, type PoolInfo } from 'poolifier'; |
8114d10e | 5 | |
268a74bb JB |
6 | import { WorkerAbstract } from './WorkerAbstract'; |
7 | import type { WorkerData, WorkerOptions } from './WorkerTypes'; | |
789871d6 | 8 | import { defaultErrorHandler, defaultExitHandler, sleep } from './WorkerUtils'; |
a4624c96 | 9 | |
268a74bb | 10 | export class WorkerDynamicPool extends WorkerAbstract<WorkerData> { |
f2bf9948 | 11 | private readonly pool: DynamicThreadPool<WorkerData>; |
a4624c96 JB |
12 | |
13 | /** | |
14 | * Create a new `WorkerDynamicPool`. | |
15 | * | |
0e4fa348 JB |
16 | * @param workerScript - |
17 | * @param workerOptions - | |
a4624c96 | 18 | */ |
4d7227e6 JB |
19 | constructor(workerScript: string, workerOptions?: WorkerOptions) { |
20 | super(workerScript, workerOptions); | |
0e4fa348 | 21 | this.workerOptions.poolOptions.errorHandler = ( |
789871d6 | 22 | this.workerOptions?.poolOptions?.errorHandler ?? defaultErrorHandler |
0e4fa348 JB |
23 | ).bind(this) as ErrorHandler<Worker>; |
24 | this.workerOptions.poolOptions.exitHandler = ( | |
789871d6 | 25 | this.workerOptions?.poolOptions?.exitHandler ?? defaultExitHandler |
0e4fa348 JB |
26 | ).bind(this) as ExitHandler<Worker>; |
27 | this.workerOptions.poolOptions.messageHandler.bind(this); | |
e7aeea18 JB |
28 | this.pool = new DynamicThreadPool<WorkerData>( |
29 | this.workerOptions.poolMinSize, | |
30 | this.workerOptions.poolMaxSize, | |
31 | this.workerScript, | |
32 | this.workerOptions.poolOptions | |
33 | ); | |
a4624c96 JB |
34 | } |
35 | ||
b779c0f8 JB |
36 | get info(): PoolInfo { |
37 | return this.pool.info; | |
38 | } | |
39 | ||
a4624c96 | 40 | get size(): number { |
b779c0f8 | 41 | return this.pool.info.workerNodes; |
a4624c96 JB |
42 | } |
43 | ||
72092cfc JB |
44 | get maxElementsPerWorker(): number | undefined { |
45 | return undefined; | |
a4624c96 JB |
46 | } |
47 | ||
962a8159 JB |
48 | get emitter(): EventEmitterAsyncResource | undefined { |
49 | return this.pool?.emitter; | |
50 | } | |
51 | ||
8baf3f8f | 52 | /** @inheritDoc */ |
6c3cfef8 JB |
53 | public async start(): Promise<void> { |
54 | // This is intentional | |
55 | } | |
a4624c96 | 56 | |
8baf3f8f | 57 | /** @inheritDoc */ |
ded13d97 JB |
58 | public async stop(): Promise<void> { |
59 | return this.pool.destroy(); | |
60 | } | |
61 | ||
8baf3f8f | 62 | /** @inheritDoc */ |
c3ee95af | 63 | public async addElement(elementData: WorkerData): Promise<void> { |
a4624c96 | 64 | await this.pool.execute(elementData); |
4bfd80fa | 65 | // Start element sequentially to optimize memory at startup |
789871d6 | 66 | this.workerOptions.elementStartDelay > 0 && (await sleep(this.workerOptions.elementStartDelay)); |
a4624c96 JB |
67 | } |
68 | } |