fix: test for worker file existence
[poolifier.git] / tests / pools / abstract / abstract-pool.test.js
CommitLineData
a61a0724 1const { expect } = require('expect')
e843b904 2const {
70a4f5ea 3 DynamicClusterPool,
9e619829 4 DynamicThreadPool,
aee46736 5 FixedClusterPool,
e843b904 6 FixedThreadPool,
aee46736 7 PoolEvents,
184855e6 8 PoolTypes,
3d6dd312 9 WorkerChoiceStrategies,
184855e6 10 WorkerTypes
cdace0e5 11} = require('../../../lib')
78099a15 12const { CircularArray } = require('../../../lib/circular-array')
29ee7e9a 13const { Queue } = require('../../../lib/queue')
23ccf9d7 14const { version } = require('../../../package.json')
2431bdb4 15const { waitPoolEvents } = require('../../test-utils')
e1ffb94f
JB
16
17describe('Abstract pool test suite', () => {
fc027381 18 const numberOfWorkers = 2
a8884ffd 19 class StubPoolWithIsMain extends FixedThreadPool {
e1ffb94f
JB
20 isMain () {
21 return false
22 }
3ec964d6 23 }
3ec964d6 24
3ec964d6 25 it('Simulate pool creation from a non main thread/process', () => {
8d3782fa
JB
26 expect(
27 () =>
a8884ffd 28 new StubPoolWithIsMain(
7c0ba920 29 numberOfWorkers,
8d3782fa
JB
30 './tests/worker-files/thread/testWorker.js',
31 {
32 errorHandler: e => console.error(e)
33 }
34 )
d4aeae5a 35 ).toThrowError('Cannot start a pool from a worker!')
3ec964d6 36 })
c510fea7
APA
37
38 it('Verify that filePath is checked', () => {
292ad316
JB
39 const expectedError = new Error(
40 'Please specify a file with a worker implementation'
41 )
7c0ba920 42 expect(() => new FixedThreadPool(numberOfWorkers)).toThrowError(
292ad316 43 expectedError
8d3782fa 44 )
7c0ba920 45 expect(() => new FixedThreadPool(numberOfWorkers, '')).toThrowError(
292ad316 46 expectedError
8d3782fa 47 )
3d6dd312
JB
48 expect(() => new FixedThreadPool(numberOfWorkers, 0)).toThrowError(
49 expectedError
50 )
51 expect(() => new FixedThreadPool(numberOfWorkers, true)).toThrowError(
52 expectedError
53 )
54 expect(
55 () => new FixedThreadPool(numberOfWorkers, './dummyWorker.ts')
56 ).toThrowError(new Error("Cannot find the worker file './dummyWorker.ts'"))
8d3782fa
JB
57 })
58
59 it('Verify that numberOfWorkers is checked', () => {
60 expect(() => new FixedThreadPool()).toThrowError(
d4aeae5a 61 'Cannot instantiate a pool without specifying the number of workers'
8d3782fa
JB
62 )
63 })
64
65 it('Verify that a negative number of workers is checked', () => {
66 expect(
67 () =>
68 new FixedClusterPool(-1, './tests/worker-files/cluster/testWorker.js')
69 ).toThrowError(
473c717a
JB
70 new RangeError(
71 'Cannot instantiate a pool with a negative number of workers'
72 )
8d3782fa
JB
73 )
74 })
75
76 it('Verify that a non integer number of workers is checked', () => {
77 expect(
78 () =>
79 new FixedThreadPool(0.25, './tests/worker-files/thread/testWorker.js')
80 ).toThrowError(
473c717a 81 new TypeError(
0d80593b 82 'Cannot instantiate a pool with a non safe integer number of workers'
8d3782fa
JB
83 )
84 )
c510fea7 85 })
7c0ba920 86
216541b6 87 it('Verify that dynamic pool sizing is checked', () => {
2431bdb4
JB
88 expect(
89 () =>
90 new DynamicThreadPool(2, 1, './tests/worker-files/thread/testWorker.js')
91 ).toThrowError(
92 new RangeError(
93 'Cannot instantiate a dynamic pool with a maximum pool size inferior to the minimum pool size'
94 )
95 )
96 expect(
97 () =>
98 new DynamicThreadPool(1, 1, './tests/worker-files/thread/testWorker.js')
99 ).toThrowError(
100 new RangeError(
101 'Cannot instantiate a dynamic pool with a minimum pool size equal to the maximum pool size. Use a fixed pool instead'
102 )
103 )
21f710aa
JB
104 expect(
105 () =>
106 new DynamicThreadPool(0, 0, './tests/worker-files/thread/testWorker.js')
107 ).toThrowError(
108 new RangeError(
109 'Cannot instantiate a dynamic pool with a minimum pool size and a maximum pool size equal to zero'
110 )
111 )
2431bdb4
JB
112 })
113
fd7ebd49 114 it('Verify that pool options are checked', async () => {
7c0ba920
JB
115 let pool = new FixedThreadPool(
116 numberOfWorkers,
117 './tests/worker-files/thread/testWorker.js'
118 )
7c0ba920 119 expect(pool.emitter).toBeDefined()
1f68cede
JB
120 expect(pool.opts.enableEvents).toBe(true)
121 expect(pool.opts.restartWorkerOnError).toBe(true)
ff733df7 122 expect(pool.opts.enableTasksQueue).toBe(false)
d4aeae5a 123 expect(pool.opts.tasksQueueOptions).toBeUndefined()
e843b904
JB
124 expect(pool.opts.workerChoiceStrategy).toBe(
125 WorkerChoiceStrategies.ROUND_ROBIN
126 )
da309861 127 expect(pool.opts.workerChoiceStrategyOptions).toStrictEqual({
932fc8be 128 runTime: { median: false },
5df69fab
JB
129 waitTime: { median: false },
130 elu: { median: false }
da309861 131 })
35cf1c03
JB
132 expect(pool.opts.messageHandler).toBeUndefined()
133 expect(pool.opts.errorHandler).toBeUndefined()
134 expect(pool.opts.onlineHandler).toBeUndefined()
135 expect(pool.opts.exitHandler).toBeUndefined()
fd7ebd49 136 await pool.destroy()
35cf1c03 137 const testHandler = () => console.log('test handler executed')
7c0ba920
JB
138 pool = new FixedThreadPool(
139 numberOfWorkers,
140 './tests/worker-files/thread/testWorker.js',
141 {
e4543b14 142 workerChoiceStrategy: WorkerChoiceStrategies.LEAST_USED,
49be33fe 143 workerChoiceStrategyOptions: {
932fc8be 144 runTime: { median: true },
fc027381 145 weights: { 0: 300, 1: 200 }
49be33fe 146 },
35cf1c03 147 enableEvents: false,
1f68cede 148 restartWorkerOnError: false,
ff733df7 149 enableTasksQueue: true,
d4aeae5a 150 tasksQueueOptions: { concurrency: 2 },
35cf1c03
JB
151 messageHandler: testHandler,
152 errorHandler: testHandler,
153 onlineHandler: testHandler,
154 exitHandler: testHandler
7c0ba920
JB
155 }
156 )
7c0ba920 157 expect(pool.emitter).toBeUndefined()
1f68cede
JB
158 expect(pool.opts.enableEvents).toBe(false)
159 expect(pool.opts.restartWorkerOnError).toBe(false)
ff733df7 160 expect(pool.opts.enableTasksQueue).toBe(true)
d4aeae5a 161 expect(pool.opts.tasksQueueOptions).toStrictEqual({ concurrency: 2 })
e843b904 162 expect(pool.opts.workerChoiceStrategy).toBe(
e4543b14 163 WorkerChoiceStrategies.LEAST_USED
e843b904 164 )
da309861 165 expect(pool.opts.workerChoiceStrategyOptions).toStrictEqual({
932fc8be 166 runTime: { median: true },
fc027381 167 weights: { 0: 300, 1: 200 }
da309861 168 })
35cf1c03
JB
169 expect(pool.opts.messageHandler).toStrictEqual(testHandler)
170 expect(pool.opts.errorHandler).toStrictEqual(testHandler)
171 expect(pool.opts.onlineHandler).toStrictEqual(testHandler)
172 expect(pool.opts.exitHandler).toStrictEqual(testHandler)
fd7ebd49 173 await pool.destroy()
7c0ba920
JB
174 })
175
a20f0ba5 176 it('Verify that pool options are validated', async () => {
d4aeae5a
JB
177 expect(
178 () =>
179 new FixedThreadPool(
180 numberOfWorkers,
181 './tests/worker-files/thread/testWorker.js',
182 {
f0d7f803 183 workerChoiceStrategy: 'invalidStrategy'
d4aeae5a
JB
184 }
185 )
f0d7f803 186 ).toThrowError("Invalid worker choice strategy 'invalidStrategy'")
d4aeae5a
JB
187 expect(
188 () =>
189 new FixedThreadPool(
190 numberOfWorkers,
191 './tests/worker-files/thread/testWorker.js',
192 {
f0d7f803 193 workerChoiceStrategyOptions: 'invalidOptions'
d4aeae5a
JB
194 }
195 )
f0d7f803
JB
196 ).toThrowError(
197 'Invalid worker choice strategy options: must be a plain object'
198 )
49be33fe
JB
199 expect(
200 () =>
201 new FixedThreadPool(
202 numberOfWorkers,
203 './tests/worker-files/thread/testWorker.js',
204 {
205 workerChoiceStrategyOptions: { weights: {} }
206 }
207 )
208 ).toThrowError(
209 'Invalid worker choice strategy options: must have a weight for each worker node'
210 )
f0d7f803
JB
211 expect(
212 () =>
213 new FixedThreadPool(
214 numberOfWorkers,
215 './tests/worker-files/thread/testWorker.js',
216 {
217 workerChoiceStrategyOptions: { measurement: 'invalidMeasurement' }
218 }
219 )
220 ).toThrowError(
221 "Invalid worker choice strategy options: invalid measurement 'invalidMeasurement'"
222 )
223 expect(
224 () =>
225 new FixedThreadPool(
226 numberOfWorkers,
227 './tests/worker-files/thread/testWorker.js',
228 {
229 enableTasksQueue: true,
230 tasksQueueOptions: { concurrency: 0 }
231 }
232 )
233 ).toThrowError("Invalid worker tasks concurrency '0'")
234 expect(
235 () =>
236 new FixedThreadPool(
237 numberOfWorkers,
238 './tests/worker-files/thread/testWorker.js',
239 {
240 enableTasksQueue: true,
241 tasksQueueOptions: 'invalidTasksQueueOptions'
242 }
243 )
244 ).toThrowError('Invalid tasks queue options: must be a plain object')
245 expect(
246 () =>
247 new FixedThreadPool(
248 numberOfWorkers,
249 './tests/worker-files/thread/testWorker.js',
250 {
251 enableTasksQueue: true,
252 tasksQueueOptions: { concurrency: 0.2 }
253 }
254 )
255 ).toThrowError('Invalid worker tasks concurrency: must be an integer')
d4aeae5a
JB
256 })
257
2431bdb4 258 it('Verify that pool worker choice strategy options can be set', async () => {
a20f0ba5
JB
259 const pool = new FixedThreadPool(
260 numberOfWorkers,
261 './tests/worker-files/thread/testWorker.js',
262 { workerChoiceStrategy: WorkerChoiceStrategies.FAIR_SHARE }
263 )
264 expect(pool.opts.workerChoiceStrategyOptions).toStrictEqual({
932fc8be 265 runTime: { median: false },
5df69fab
JB
266 waitTime: { median: false },
267 elu: { median: false }
a20f0ba5
JB
268 })
269 for (const [, workerChoiceStrategy] of pool.workerChoiceStrategyContext
270 .workerChoiceStrategies) {
86bf340d 271 expect(workerChoiceStrategy.opts).toStrictEqual({
932fc8be 272 runTime: { median: false },
5df69fab
JB
273 waitTime: { median: false },
274 elu: { median: false }
86bf340d 275 })
a20f0ba5 276 }
87de9ff5
JB
277 expect(
278 pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
279 ).toStrictEqual({
932fc8be
JB
280 runTime: {
281 aggregate: true,
282 average: true,
283 median: false
284 },
285 waitTime: {
286 aggregate: false,
287 average: false,
288 median: false
289 },
5df69fab 290 elu: {
9adcefab
JB
291 aggregate: true,
292 average: true,
5df69fab
JB
293 median: false
294 }
86bf340d 295 })
9adcefab
JB
296 pool.setWorkerChoiceStrategyOptions({
297 runTime: { median: true },
298 elu: { median: true }
299 })
a20f0ba5 300 expect(pool.opts.workerChoiceStrategyOptions).toStrictEqual({
9adcefab
JB
301 runTime: { median: true },
302 elu: { median: true }
a20f0ba5
JB
303 })
304 for (const [, workerChoiceStrategy] of pool.workerChoiceStrategyContext
305 .workerChoiceStrategies) {
932fc8be 306 expect(workerChoiceStrategy.opts).toStrictEqual({
9adcefab
JB
307 runTime: { median: true },
308 elu: { median: true }
932fc8be 309 })
a20f0ba5 310 }
87de9ff5
JB
311 expect(
312 pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
313 ).toStrictEqual({
932fc8be
JB
314 runTime: {
315 aggregate: true,
316 average: false,
317 median: true
318 },
319 waitTime: {
320 aggregate: false,
321 average: false,
322 median: false
323 },
5df69fab 324 elu: {
9adcefab 325 aggregate: true,
5df69fab 326 average: false,
9adcefab 327 median: true
5df69fab 328 }
86bf340d 329 })
9adcefab
JB
330 pool.setWorkerChoiceStrategyOptions({
331 runTime: { median: false },
332 elu: { median: false }
333 })
a20f0ba5 334 expect(pool.opts.workerChoiceStrategyOptions).toStrictEqual({
9adcefab
JB
335 runTime: { median: false },
336 elu: { median: false }
a20f0ba5
JB
337 })
338 for (const [, workerChoiceStrategy] of pool.workerChoiceStrategyContext
339 .workerChoiceStrategies) {
932fc8be 340 expect(workerChoiceStrategy.opts).toStrictEqual({
9adcefab
JB
341 runTime: { median: false },
342 elu: { median: false }
932fc8be 343 })
a20f0ba5 344 }
87de9ff5
JB
345 expect(
346 pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
347 ).toStrictEqual({
932fc8be
JB
348 runTime: {
349 aggregate: true,
350 average: true,
351 median: false
352 },
353 waitTime: {
354 aggregate: false,
355 average: false,
356 median: false
357 },
5df69fab 358 elu: {
9adcefab
JB
359 aggregate: true,
360 average: true,
5df69fab
JB
361 median: false
362 }
86bf340d 363 })
1f95d544
JB
364 expect(() =>
365 pool.setWorkerChoiceStrategyOptions('invalidWorkerChoiceStrategyOptions')
366 ).toThrowError(
367 'Invalid worker choice strategy options: must be a plain object'
368 )
369 expect(() =>
370 pool.setWorkerChoiceStrategyOptions({ weights: {} })
371 ).toThrowError(
372 'Invalid worker choice strategy options: must have a weight for each worker node'
373 )
374 expect(() =>
375 pool.setWorkerChoiceStrategyOptions({ measurement: 'invalidMeasurement' })
376 ).toThrowError(
377 "Invalid worker choice strategy options: invalid measurement 'invalidMeasurement'"
378 )
a20f0ba5
JB
379 await pool.destroy()
380 })
381
2431bdb4 382 it('Verify that pool tasks queue can be enabled/disabled', async () => {
a20f0ba5
JB
383 const pool = new FixedThreadPool(
384 numberOfWorkers,
385 './tests/worker-files/thread/testWorker.js'
386 )
387 expect(pool.opts.enableTasksQueue).toBe(false)
388 expect(pool.opts.tasksQueueOptions).toBeUndefined()
389 pool.enableTasksQueue(true)
390 expect(pool.opts.enableTasksQueue).toBe(true)
391 expect(pool.opts.tasksQueueOptions).toStrictEqual({ concurrency: 1 })
392 pool.enableTasksQueue(true, { concurrency: 2 })
393 expect(pool.opts.enableTasksQueue).toBe(true)
394 expect(pool.opts.tasksQueueOptions).toStrictEqual({ concurrency: 2 })
395 pool.enableTasksQueue(false)
396 expect(pool.opts.enableTasksQueue).toBe(false)
397 expect(pool.opts.tasksQueueOptions).toBeUndefined()
398 await pool.destroy()
399 })
400
2431bdb4 401 it('Verify that pool tasks queue options can be set', async () => {
a20f0ba5
JB
402 const pool = new FixedThreadPool(
403 numberOfWorkers,
404 './tests/worker-files/thread/testWorker.js',
405 { enableTasksQueue: true }
406 )
407 expect(pool.opts.tasksQueueOptions).toStrictEqual({ concurrency: 1 })
408 pool.setTasksQueueOptions({ concurrency: 2 })
409 expect(pool.opts.tasksQueueOptions).toStrictEqual({ concurrency: 2 })
f0d7f803
JB
410 expect(() =>
411 pool.setTasksQueueOptions('invalidTasksQueueOptions')
412 ).toThrowError('Invalid tasks queue options: must be a plain object')
a20f0ba5
JB
413 expect(() => pool.setTasksQueueOptions({ concurrency: 0 })).toThrowError(
414 "Invalid worker tasks concurrency '0'"
415 )
f0d7f803
JB
416 expect(() => pool.setTasksQueueOptions({ concurrency: 0.2 })).toThrowError(
417 'Invalid worker tasks concurrency: must be an integer'
418 )
a20f0ba5
JB
419 await pool.destroy()
420 })
421
6b27d407
JB
422 it('Verify that pool info is set', async () => {
423 let pool = new FixedThreadPool(
424 numberOfWorkers,
425 './tests/worker-files/thread/testWorker.js'
426 )
427 expect(pool.info).toStrictEqual({
23ccf9d7 428 version,
6b27d407 429 type: PoolTypes.fixed,
184855e6 430 worker: WorkerTypes.thread,
2431bdb4
JB
431 ready: false,
432 strategy: WorkerChoiceStrategies.ROUND_ROBIN,
6b27d407
JB
433 minSize: numberOfWorkers,
434 maxSize: numberOfWorkers,
435 workerNodes: numberOfWorkers,
436 idleWorkerNodes: numberOfWorkers,
437 busyWorkerNodes: 0,
a4e07f72
JB
438 executedTasks: 0,
439 executingTasks: 0,
6b27d407 440 queuedTasks: 0,
a4e07f72
JB
441 maxQueuedTasks: 0,
442 failedTasks: 0
6b27d407 443 })
2dca6cab
JB
444 await waitPoolEvents(pool, PoolEvents.ready, 1)
445 expect(pool.info).toStrictEqual({
446 version,
447 type: PoolTypes.fixed,
448 worker: WorkerTypes.thread,
449 ready: true,
450 strategy: WorkerChoiceStrategies.ROUND_ROBIN,
451 minSize: numberOfWorkers,
452 maxSize: numberOfWorkers,
453 workerNodes: numberOfWorkers,
454 idleWorkerNodes: numberOfWorkers,
455 busyWorkerNodes: 0,
456 executedTasks: 0,
457 executingTasks: 0,
458 queuedTasks: 0,
459 maxQueuedTasks: 0,
460 failedTasks: 0
461 })
6b27d407
JB
462 await pool.destroy()
463 pool = new DynamicClusterPool(
2431bdb4 464 Math.floor(numberOfWorkers / 2),
6b27d407 465 numberOfWorkers,
ecdfbdc0 466 './tests/worker-files/cluster/testWorker.js'
6b27d407
JB
467 )
468 expect(pool.info).toStrictEqual({
23ccf9d7 469 version,
6b27d407 470 type: PoolTypes.dynamic,
184855e6 471 worker: WorkerTypes.cluster,
2431bdb4
JB
472 ready: false,
473 strategy: WorkerChoiceStrategies.ROUND_ROBIN,
474 minSize: Math.floor(numberOfWorkers / 2),
475 maxSize: numberOfWorkers,
476 workerNodes: Math.floor(numberOfWorkers / 2),
477 idleWorkerNodes: Math.floor(numberOfWorkers / 2),
6b27d407 478 busyWorkerNodes: 0,
a4e07f72
JB
479 executedTasks: 0,
480 executingTasks: 0,
6b27d407 481 queuedTasks: 0,
a4e07f72
JB
482 maxQueuedTasks: 0,
483 failedTasks: 0
6b27d407 484 })
2dca6cab
JB
485 await waitPoolEvents(pool, PoolEvents.ready, 1)
486 expect(pool.info).toStrictEqual({
487 version,
488 type: PoolTypes.dynamic,
489 worker: WorkerTypes.cluster,
490 ready: true,
491 strategy: WorkerChoiceStrategies.ROUND_ROBIN,
492 minSize: Math.floor(numberOfWorkers / 2),
493 maxSize: numberOfWorkers,
494 workerNodes: Math.floor(numberOfWorkers / 2),
495 idleWorkerNodes: Math.floor(numberOfWorkers / 2),
496 busyWorkerNodes: 0,
497 executedTasks: 0,
498 executingTasks: 0,
499 queuedTasks: 0,
500 maxQueuedTasks: 0,
501 failedTasks: 0
502 })
6b27d407
JB
503 await pool.destroy()
504 })
505
2431bdb4 506 it('Verify that pool worker tasks usage are initialized', async () => {
bf9549ae
JB
507 const pool = new FixedClusterPool(
508 numberOfWorkers,
509 './tests/worker-files/cluster/testWorker.js'
510 )
f06e48d8 511 for (const workerNode of pool.workerNodes) {
465b2940 512 expect(workerNode.usage).toStrictEqual({
a4e07f72
JB
513 tasks: {
514 executed: 0,
515 executing: 0,
516 queued: 0,
df593701 517 maxQueued: 0,
a4e07f72
JB
518 failed: 0
519 },
520 runTime: {
a4e07f72
JB
521 history: expect.any(CircularArray)
522 },
523 waitTime: {
a4e07f72
JB
524 history: expect.any(CircularArray)
525 },
5df69fab
JB
526 elu: {
527 idle: {
5df69fab
JB
528 history: expect.any(CircularArray)
529 },
530 active: {
5df69fab 531 history: expect.any(CircularArray)
f7510105 532 }
5df69fab 533 }
86bf340d 534 })
f06e48d8
JB
535 }
536 await pool.destroy()
537 })
538
2431bdb4
JB
539 it('Verify that pool worker tasks queue are initialized', async () => {
540 let pool = new FixedClusterPool(
f06e48d8
JB
541 numberOfWorkers,
542 './tests/worker-files/cluster/testWorker.js'
543 )
544 for (const workerNode of pool.workerNodes) {
545 expect(workerNode.tasksQueue).toBeDefined()
29ee7e9a 546 expect(workerNode.tasksQueue).toBeInstanceOf(Queue)
4d8bf9e4 547 expect(workerNode.tasksQueue.size).toBe(0)
9c16fb4b 548 expect(workerNode.tasksQueue.maxSize).toBe(0)
bf9549ae 549 }
fd7ebd49 550 await pool.destroy()
2431bdb4
JB
551 pool = new DynamicThreadPool(
552 Math.floor(numberOfWorkers / 2),
553 numberOfWorkers,
554 './tests/worker-files/thread/testWorker.js'
555 )
556 for (const workerNode of pool.workerNodes) {
557 expect(workerNode.tasksQueue).toBeDefined()
558 expect(workerNode.tasksQueue).toBeInstanceOf(Queue)
559 expect(workerNode.tasksQueue.size).toBe(0)
560 expect(workerNode.tasksQueue.maxSize).toBe(0)
561 }
562 })
563
564 it('Verify that pool worker info are initialized', async () => {
565 let pool = new FixedClusterPool(
566 numberOfWorkers,
567 './tests/worker-files/cluster/testWorker.js'
568 )
569 for (const workerNode of pool.workerNodes) {
570 expect(workerNode.info).toStrictEqual({
571 id: expect.any(Number),
572 type: WorkerTypes.cluster,
573 dynamic: false,
574 ready: false
575 })
576 }
2dca6cab
JB
577 await waitPoolEvents(pool, PoolEvents.ready, 1)
578 for (const workerNode of pool.workerNodes) {
579 expect(workerNode.info).toStrictEqual({
580 id: expect.any(Number),
581 type: WorkerTypes.cluster,
582 dynamic: false,
583 ready: true
584 })
585 }
2431bdb4
JB
586 await pool.destroy()
587 pool = new DynamicThreadPool(
588 Math.floor(numberOfWorkers / 2),
589 numberOfWorkers,
590 './tests/worker-files/thread/testWorker.js'
591 )
592 for (const workerNode of pool.workerNodes) {
593 expect(workerNode.info).toStrictEqual({
594 id: expect.any(Number),
595 type: WorkerTypes.thread,
596 dynamic: false,
597 ready: false
598 })
599 }
2dca6cab
JB
600 await waitPoolEvents(pool, PoolEvents.ready, 1)
601 for (const workerNode of pool.workerNodes) {
602 expect(workerNode.info).toStrictEqual({
603 id: expect.any(Number),
604 type: WorkerTypes.thread,
605 dynamic: false,
606 ready: true
607 })
608 }
bf9549ae
JB
609 })
610
2431bdb4 611 it('Verify that pool worker tasks usage are computed', async () => {
bf9549ae
JB
612 const pool = new FixedClusterPool(
613 numberOfWorkers,
614 './tests/worker-files/cluster/testWorker.js'
615 )
09c2d0d3 616 const promises = new Set()
fc027381
JB
617 const maxMultiplier = 2
618 for (let i = 0; i < numberOfWorkers * maxMultiplier; i++) {
09c2d0d3 619 promises.add(pool.execute())
bf9549ae 620 }
f06e48d8 621 for (const workerNode of pool.workerNodes) {
465b2940 622 expect(workerNode.usage).toStrictEqual({
a4e07f72
JB
623 tasks: {
624 executed: 0,
625 executing: maxMultiplier,
626 queued: 0,
df593701 627 maxQueued: 0,
a4e07f72
JB
628 failed: 0
629 },
630 runTime: {
a4e07f72
JB
631 history: expect.any(CircularArray)
632 },
633 waitTime: {
a4e07f72
JB
634 history: expect.any(CircularArray)
635 },
5df69fab
JB
636 elu: {
637 idle: {
5df69fab
JB
638 history: expect.any(CircularArray)
639 },
640 active: {
5df69fab 641 history: expect.any(CircularArray)
f7510105 642 }
5df69fab 643 }
86bf340d 644 })
bf9549ae
JB
645 }
646 await Promise.all(promises)
f06e48d8 647 for (const workerNode of pool.workerNodes) {
465b2940 648 expect(workerNode.usage).toStrictEqual({
a4e07f72
JB
649 tasks: {
650 executed: maxMultiplier,
651 executing: 0,
652 queued: 0,
df593701 653 maxQueued: 0,
a4e07f72
JB
654 failed: 0
655 },
656 runTime: {
a4e07f72
JB
657 history: expect.any(CircularArray)
658 },
659 waitTime: {
a4e07f72
JB
660 history: expect.any(CircularArray)
661 },
5df69fab
JB
662 elu: {
663 idle: {
5df69fab
JB
664 history: expect.any(CircularArray)
665 },
666 active: {
5df69fab 667 history: expect.any(CircularArray)
f7510105 668 }
5df69fab 669 }
86bf340d 670 })
bf9549ae 671 }
fd7ebd49 672 await pool.destroy()
bf9549ae
JB
673 })
674
2431bdb4 675 it('Verify that pool worker tasks usage are reset at worker choice strategy change', async () => {
7fd82a1c 676 const pool = new DynamicThreadPool(
2431bdb4 677 Math.floor(numberOfWorkers / 2),
8f4878b7 678 numberOfWorkers,
9e619829
JB
679 './tests/worker-files/thread/testWorker.js'
680 )
09c2d0d3 681 const promises = new Set()
ee9f5295
JB
682 const maxMultiplier = 2
683 for (let i = 0; i < numberOfWorkers * maxMultiplier; i++) {
09c2d0d3 684 promises.add(pool.execute())
9e619829
JB
685 }
686 await Promise.all(promises)
f06e48d8 687 for (const workerNode of pool.workerNodes) {
465b2940 688 expect(workerNode.usage).toStrictEqual({
a4e07f72
JB
689 tasks: {
690 executed: expect.any(Number),
691 executing: 0,
692 queued: 0,
df593701 693 maxQueued: 0,
a4e07f72
JB
694 failed: 0
695 },
696 runTime: {
a4e07f72
JB
697 history: expect.any(CircularArray)
698 },
699 waitTime: {
a4e07f72
JB
700 history: expect.any(CircularArray)
701 },
5df69fab
JB
702 elu: {
703 idle: {
5df69fab
JB
704 history: expect.any(CircularArray)
705 },
706 active: {
5df69fab 707 history: expect.any(CircularArray)
f7510105 708 }
5df69fab 709 }
86bf340d 710 })
465b2940
JB
711 expect(workerNode.usage.tasks.executed).toBeGreaterThan(0)
712 expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(maxMultiplier)
9e619829
JB
713 }
714 pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.FAIR_SHARE)
f06e48d8 715 for (const workerNode of pool.workerNodes) {
465b2940 716 expect(workerNode.usage).toStrictEqual({
a4e07f72
JB
717 tasks: {
718 executed: 0,
719 executing: 0,
720 queued: 0,
df593701 721 maxQueued: 0,
a4e07f72
JB
722 failed: 0
723 },
724 runTime: {
a4e07f72
JB
725 history: expect.any(CircularArray)
726 },
727 waitTime: {
a4e07f72
JB
728 history: expect.any(CircularArray)
729 },
5df69fab
JB
730 elu: {
731 idle: {
5df69fab
JB
732 history: expect.any(CircularArray)
733 },
734 active: {
5df69fab 735 history: expect.any(CircularArray)
f7510105 736 }
5df69fab 737 }
86bf340d 738 })
465b2940
JB
739 expect(workerNode.usage.runTime.history.length).toBe(0)
740 expect(workerNode.usage.waitTime.history.length).toBe(0)
ee11a4a2 741 }
fd7ebd49 742 await pool.destroy()
ee11a4a2
JB
743 })
744
164d950a
JB
745 it("Verify that pool event emitter 'full' event can register a callback", async () => {
746 const pool = new DynamicThreadPool(
2431bdb4 747 Math.floor(numberOfWorkers / 2),
164d950a
JB
748 numberOfWorkers,
749 './tests/worker-files/thread/testWorker.js'
750 )
09c2d0d3 751 const promises = new Set()
164d950a 752 let poolFull = 0
d46660cd
JB
753 let poolInfo
754 pool.emitter.on(PoolEvents.full, info => {
755 ++poolFull
756 poolInfo = info
757 })
164d950a 758 for (let i = 0; i < numberOfWorkers * 2; i++) {
f5d14e90 759 promises.add(pool.execute())
164d950a
JB
760 }
761 await Promise.all(promises)
2431bdb4
JB
762 // The `full` event is triggered when the number of submitted tasks at once reach the maximum number of workers in the dynamic pool.
763 // So in total numberOfWorkers * 2 - 1 times for a loop submitting up to numberOfWorkers * 2 tasks to the dynamic pool with min = (max = numberOfWorkers) / 2.
764 expect(poolFull).toBe(numberOfWorkers * 2 - 1)
d46660cd 765 expect(poolInfo).toStrictEqual({
23ccf9d7 766 version,
d46660cd
JB
767 type: PoolTypes.dynamic,
768 worker: WorkerTypes.thread,
2431bdb4
JB
769 ready: expect.any(Boolean),
770 strategy: WorkerChoiceStrategies.ROUND_ROBIN,
771 minSize: expect.any(Number),
772 maxSize: expect.any(Number),
773 workerNodes: expect.any(Number),
774 idleWorkerNodes: expect.any(Number),
775 busyWorkerNodes: expect.any(Number),
776 executedTasks: expect.any(Number),
777 executingTasks: expect.any(Number),
778 queuedTasks: expect.any(Number),
779 maxQueuedTasks: expect.any(Number),
780 failedTasks: expect.any(Number)
781 })
782 await pool.destroy()
783 })
784
785 it("Verify that pool event emitter 'ready' event can register a callback", async () => {
d5024c00
JB
786 const pool = new DynamicClusterPool(
787 Math.floor(numberOfWorkers / 2),
2431bdb4
JB
788 numberOfWorkers,
789 './tests/worker-files/cluster/testWorker.js'
790 )
791 let poolReady = 0
792 let poolInfo
793 pool.emitter.on(PoolEvents.ready, info => {
794 ++poolReady
795 poolInfo = info
796 })
797 await waitPoolEvents(pool, PoolEvents.ready, 1)
798 expect(poolReady).toBe(1)
799 expect(poolInfo).toStrictEqual({
800 version,
d5024c00 801 type: PoolTypes.dynamic,
2431bdb4
JB
802 worker: WorkerTypes.cluster,
803 ready: true,
804 strategy: WorkerChoiceStrategies.ROUND_ROBIN,
d46660cd
JB
805 minSize: expect.any(Number),
806 maxSize: expect.any(Number),
807 workerNodes: expect.any(Number),
808 idleWorkerNodes: expect.any(Number),
809 busyWorkerNodes: expect.any(Number),
a4e07f72
JB
810 executedTasks: expect.any(Number),
811 executingTasks: expect.any(Number),
d46660cd 812 queuedTasks: expect.any(Number),
a4e07f72
JB
813 maxQueuedTasks: expect.any(Number),
814 failedTasks: expect.any(Number)
d46660cd 815 })
164d950a
JB
816 await pool.destroy()
817 })
818
cf597bc5 819 it("Verify that pool event emitter 'busy' event can register a callback", async () => {
7c0ba920
JB
820 const pool = new FixedThreadPool(
821 numberOfWorkers,
822 './tests/worker-files/thread/testWorker.js'
823 )
09c2d0d3 824 const promises = new Set()
7c0ba920 825 let poolBusy = 0
d46660cd
JB
826 let poolInfo
827 pool.emitter.on(PoolEvents.busy, info => {
828 ++poolBusy
829 poolInfo = info
830 })
7c0ba920 831 for (let i = 0; i < numberOfWorkers * 2; i++) {
f5d14e90 832 promises.add(pool.execute())
7c0ba920 833 }
cf597bc5 834 await Promise.all(promises)
14916bf9
JB
835 // The `busy` event is triggered when the number of submitted tasks at once reach the number of fixed pool workers.
836 // So in total numberOfWorkers + 1 times for a loop submitting up to numberOfWorkers * 2 tasks to the fixed pool.
837 expect(poolBusy).toBe(numberOfWorkers + 1)
d46660cd 838 expect(poolInfo).toStrictEqual({
23ccf9d7 839 version,
d46660cd
JB
840 type: PoolTypes.fixed,
841 worker: WorkerTypes.thread,
2431bdb4
JB
842 ready: expect.any(Boolean),
843 strategy: WorkerChoiceStrategies.ROUND_ROBIN,
d46660cd
JB
844 minSize: expect.any(Number),
845 maxSize: expect.any(Number),
846 workerNodes: expect.any(Number),
847 idleWorkerNodes: expect.any(Number),
848 busyWorkerNodes: expect.any(Number),
a4e07f72
JB
849 executedTasks: expect.any(Number),
850 executingTasks: expect.any(Number),
d46660cd 851 queuedTasks: expect.any(Number),
a4e07f72
JB
852 maxQueuedTasks: expect.any(Number),
853 failedTasks: expect.any(Number)
d46660cd 854 })
fd7ebd49 855 await pool.destroy()
7c0ba920 856 })
70a4f5ea
JB
857
858 it('Verify that multiple tasks worker is working', async () => {
859 const pool = new DynamicClusterPool(
2431bdb4 860 Math.floor(numberOfWorkers / 2),
70a4f5ea 861 numberOfWorkers,
70a4f5ea
JB
862 './tests/worker-files/cluster/testMultiTasksWorker.js'
863 )
864 const data = { n: 10 }
82888165 865 const result0 = await pool.execute(data)
30b963d4 866 expect(result0).toStrictEqual({ ok: 1 })
70a4f5ea 867 const result1 = await pool.execute(data, 'jsonIntegerSerialization')
30b963d4 868 expect(result1).toStrictEqual({ ok: 1 })
70a4f5ea
JB
869 const result2 = await pool.execute(data, 'factorial')
870 expect(result2).toBe(3628800)
871 const result3 = await pool.execute(data, 'fibonacci')
024daf59 872 expect(result3).toBe(55)
70a4f5ea 873 })
3ec964d6 874})