2c8d23184d4dfb49342cd4998670a1f0da1f7619
[poolifier.git] / tests / pools / abstract / abstract-pool.test.js
1 const { expect } = require('expect')
2 const {
3 DynamicClusterPool,
4 DynamicThreadPool,
5 FixedClusterPool,
6 FixedThreadPool,
7 PoolEvents,
8 PoolTypes,
9 WorkerChoiceStrategies,
10 WorkerTypes
11 } = require('../../../lib')
12 const { CircularArray } = require('../../../lib/circular-array')
13 const { Queue } = require('../../../lib/queue')
14 const { version } = require('../../../package.json')
15 const { waitPoolEvents } = require('../../test-utils')
16
17 describe('Abstract pool test suite', () => {
18 const numberOfWorkers = 2
19 class StubPoolWithIsMain extends FixedThreadPool {
20 isMain () {
21 return false
22 }
23 }
24
25 it('Simulate pool creation from a non main thread/process', () => {
26 expect(
27 () =>
28 new StubPoolWithIsMain(
29 numberOfWorkers,
30 './tests/worker-files/thread/testWorker.js',
31 {
32 errorHandler: e => console.error(e)
33 }
34 )
35 ).toThrowError('Cannot start a pool from a worker!')
36 })
37
38 it('Verify that filePath is checked', () => {
39 const expectedError = new Error(
40 'Please specify a file with a worker implementation'
41 )
42 expect(() => new FixedThreadPool(numberOfWorkers)).toThrowError(
43 expectedError
44 )
45 expect(() => new FixedThreadPool(numberOfWorkers, '')).toThrowError(
46 expectedError
47 )
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'"))
57 })
58
59 it('Verify that numberOfWorkers is checked', () => {
60 expect(() => new FixedThreadPool()).toThrowError(
61 'Cannot instantiate a pool without specifying the number of workers'
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(
70 new RangeError(
71 'Cannot instantiate a pool with a negative number of workers'
72 )
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(
81 new TypeError(
82 'Cannot instantiate a pool with a non safe integer number of workers'
83 )
84 )
85 })
86
87 it('Verify that dynamic pool sizing is checked', () => {
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 )
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 pool size equal to zero'
110 )
111 )
112 })
113
114 it('Verify that pool options are checked', async () => {
115 let pool = new FixedThreadPool(
116 numberOfWorkers,
117 './tests/worker-files/thread/testWorker.js'
118 )
119 expect(pool.emitter).toBeDefined()
120 expect(pool.opts.enableEvents).toBe(true)
121 expect(pool.opts.restartWorkerOnError).toBe(true)
122 expect(pool.opts.enableTasksQueue).toBe(false)
123 expect(pool.opts.tasksQueueOptions).toBeUndefined()
124 expect(pool.opts.workerChoiceStrategy).toBe(
125 WorkerChoiceStrategies.ROUND_ROBIN
126 )
127 expect(pool.opts.workerChoiceStrategyOptions).toStrictEqual({
128 runTime: { median: false },
129 waitTime: { median: false },
130 elu: { median: false }
131 })
132 expect(pool.opts.messageHandler).toBeUndefined()
133 expect(pool.opts.errorHandler).toBeUndefined()
134 expect(pool.opts.onlineHandler).toBeUndefined()
135 expect(pool.opts.exitHandler).toBeUndefined()
136 await pool.destroy()
137 const testHandler = () => console.log('test handler executed')
138 pool = new FixedThreadPool(
139 numberOfWorkers,
140 './tests/worker-files/thread/testWorker.js',
141 {
142 workerChoiceStrategy: WorkerChoiceStrategies.LEAST_USED,
143 workerChoiceStrategyOptions: {
144 runTime: { median: true },
145 weights: { 0: 300, 1: 200 }
146 },
147 enableEvents: false,
148 restartWorkerOnError: false,
149 enableTasksQueue: true,
150 tasksQueueOptions: { concurrency: 2 },
151 messageHandler: testHandler,
152 errorHandler: testHandler,
153 onlineHandler: testHandler,
154 exitHandler: testHandler
155 }
156 )
157 expect(pool.emitter).toBeUndefined()
158 expect(pool.opts.enableEvents).toBe(false)
159 expect(pool.opts.restartWorkerOnError).toBe(false)
160 expect(pool.opts.enableTasksQueue).toBe(true)
161 expect(pool.opts.tasksQueueOptions).toStrictEqual({ concurrency: 2 })
162 expect(pool.opts.workerChoiceStrategy).toBe(
163 WorkerChoiceStrategies.LEAST_USED
164 )
165 expect(pool.opts.workerChoiceStrategyOptions).toStrictEqual({
166 runTime: { median: true },
167 weights: { 0: 300, 1: 200 }
168 })
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)
173 await pool.destroy()
174 })
175
176 it('Verify that pool options are validated', async () => {
177 expect(
178 () =>
179 new FixedThreadPool(
180 numberOfWorkers,
181 './tests/worker-files/thread/testWorker.js',
182 {
183 workerChoiceStrategy: 'invalidStrategy'
184 }
185 )
186 ).toThrowError("Invalid worker choice strategy 'invalidStrategy'")
187 expect(
188 () =>
189 new FixedThreadPool(
190 numberOfWorkers,
191 './tests/worker-files/thread/testWorker.js',
192 {
193 workerChoiceStrategyOptions: 'invalidOptions'
194 }
195 )
196 ).toThrowError(
197 'Invalid worker choice strategy options: must be a plain object'
198 )
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 )
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')
256 })
257
258 it('Verify that pool worker choice strategy options can be set', async () => {
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({
265 runTime: { median: false },
266 waitTime: { median: false },
267 elu: { median: false }
268 })
269 for (const [, workerChoiceStrategy] of pool.workerChoiceStrategyContext
270 .workerChoiceStrategies) {
271 expect(workerChoiceStrategy.opts).toStrictEqual({
272 runTime: { median: false },
273 waitTime: { median: false },
274 elu: { median: false }
275 })
276 }
277 expect(
278 pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
279 ).toStrictEqual({
280 runTime: {
281 aggregate: true,
282 average: true,
283 median: false
284 },
285 waitTime: {
286 aggregate: false,
287 average: false,
288 median: false
289 },
290 elu: {
291 aggregate: true,
292 average: true,
293 median: false
294 }
295 })
296 pool.setWorkerChoiceStrategyOptions({
297 runTime: { median: true },
298 elu: { median: true }
299 })
300 expect(pool.opts.workerChoiceStrategyOptions).toStrictEqual({
301 runTime: { median: true },
302 elu: { median: true }
303 })
304 for (const [, workerChoiceStrategy] of pool.workerChoiceStrategyContext
305 .workerChoiceStrategies) {
306 expect(workerChoiceStrategy.opts).toStrictEqual({
307 runTime: { median: true },
308 elu: { median: true }
309 })
310 }
311 expect(
312 pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
313 ).toStrictEqual({
314 runTime: {
315 aggregate: true,
316 average: false,
317 median: true
318 },
319 waitTime: {
320 aggregate: false,
321 average: false,
322 median: false
323 },
324 elu: {
325 aggregate: true,
326 average: false,
327 median: true
328 }
329 })
330 pool.setWorkerChoiceStrategyOptions({
331 runTime: { median: false },
332 elu: { median: false }
333 })
334 expect(pool.opts.workerChoiceStrategyOptions).toStrictEqual({
335 runTime: { median: false },
336 elu: { median: false }
337 })
338 for (const [, workerChoiceStrategy] of pool.workerChoiceStrategyContext
339 .workerChoiceStrategies) {
340 expect(workerChoiceStrategy.opts).toStrictEqual({
341 runTime: { median: false },
342 elu: { median: false }
343 })
344 }
345 expect(
346 pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
347 ).toStrictEqual({
348 runTime: {
349 aggregate: true,
350 average: true,
351 median: false
352 },
353 waitTime: {
354 aggregate: false,
355 average: false,
356 median: false
357 },
358 elu: {
359 aggregate: true,
360 average: true,
361 median: false
362 }
363 })
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 )
379 await pool.destroy()
380 })
381
382 it('Verify that pool tasks queue can be enabled/disabled', async () => {
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
401 it('Verify that pool tasks queue options can be set', async () => {
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 })
410 expect(() =>
411 pool.setTasksQueueOptions('invalidTasksQueueOptions')
412 ).toThrowError('Invalid tasks queue options: must be a plain object')
413 expect(() => pool.setTasksQueueOptions({ concurrency: 0 })).toThrowError(
414 "Invalid worker tasks concurrency '0'"
415 )
416 expect(() => pool.setTasksQueueOptions({ concurrency: 0.2 })).toThrowError(
417 'Invalid worker tasks concurrency: must be an integer'
418 )
419 await pool.destroy()
420 })
421
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({
428 version,
429 type: PoolTypes.fixed,
430 worker: WorkerTypes.thread,
431 ready: true,
432 strategy: WorkerChoiceStrategies.ROUND_ROBIN,
433 minSize: numberOfWorkers,
434 maxSize: numberOfWorkers,
435 workerNodes: numberOfWorkers,
436 idleWorkerNodes: numberOfWorkers,
437 busyWorkerNodes: 0,
438 executedTasks: 0,
439 executingTasks: 0,
440 queuedTasks: 0,
441 maxQueuedTasks: 0,
442 failedTasks: 0
443 })
444 await pool.destroy()
445 pool = new DynamicClusterPool(
446 Math.floor(numberOfWorkers / 2),
447 numberOfWorkers,
448 './tests/worker-files/cluster/testWorker.js'
449 )
450 expect(pool.info).toStrictEqual({
451 version,
452 type: PoolTypes.dynamic,
453 worker: WorkerTypes.cluster,
454 ready: true,
455 strategy: WorkerChoiceStrategies.ROUND_ROBIN,
456 minSize: Math.floor(numberOfWorkers / 2),
457 maxSize: numberOfWorkers,
458 workerNodes: Math.floor(numberOfWorkers / 2),
459 idleWorkerNodes: Math.floor(numberOfWorkers / 2),
460 busyWorkerNodes: 0,
461 executedTasks: 0,
462 executingTasks: 0,
463 queuedTasks: 0,
464 maxQueuedTasks: 0,
465 failedTasks: 0
466 })
467 await pool.destroy()
468 })
469
470 it('Verify that pool worker tasks usage are initialized', async () => {
471 const pool = new FixedClusterPool(
472 numberOfWorkers,
473 './tests/worker-files/cluster/testWorker.js'
474 )
475 for (const workerNode of pool.workerNodes) {
476 expect(workerNode.usage).toStrictEqual({
477 tasks: {
478 executed: 0,
479 executing: 0,
480 queued: 0,
481 maxQueued: 0,
482 failed: 0
483 },
484 runTime: {
485 history: expect.any(CircularArray)
486 },
487 waitTime: {
488 history: expect.any(CircularArray)
489 },
490 elu: {
491 idle: {
492 history: expect.any(CircularArray)
493 },
494 active: {
495 history: expect.any(CircularArray)
496 }
497 }
498 })
499 }
500 await pool.destroy()
501 })
502
503 it('Verify that pool worker tasks queue are initialized', async () => {
504 let pool = new FixedClusterPool(
505 numberOfWorkers,
506 './tests/worker-files/cluster/testWorker.js'
507 )
508 for (const workerNode of pool.workerNodes) {
509 expect(workerNode.tasksQueue).toBeDefined()
510 expect(workerNode.tasksQueue).toBeInstanceOf(Queue)
511 expect(workerNode.tasksQueue.size).toBe(0)
512 expect(workerNode.tasksQueue.maxSize).toBe(0)
513 }
514 await pool.destroy()
515 pool = new DynamicThreadPool(
516 Math.floor(numberOfWorkers / 2),
517 numberOfWorkers,
518 './tests/worker-files/thread/testWorker.js'
519 )
520 for (const workerNode of pool.workerNodes) {
521 expect(workerNode.tasksQueue).toBeDefined()
522 expect(workerNode.tasksQueue).toBeInstanceOf(Queue)
523 expect(workerNode.tasksQueue.size).toBe(0)
524 expect(workerNode.tasksQueue.maxSize).toBe(0)
525 }
526 })
527
528 it('Verify that pool worker info are initialized', async () => {
529 let pool = new FixedClusterPool(
530 numberOfWorkers,
531 './tests/worker-files/cluster/testWorker.js'
532 )
533 for (const workerNode of pool.workerNodes) {
534 expect(workerNode.info).toStrictEqual({
535 id: expect.any(Number),
536 type: WorkerTypes.cluster,
537 dynamic: false,
538 ready: true
539 })
540 }
541 await pool.destroy()
542 pool = new DynamicThreadPool(
543 Math.floor(numberOfWorkers / 2),
544 numberOfWorkers,
545 './tests/worker-files/thread/testWorker.js'
546 )
547 for (const workerNode of pool.workerNodes) {
548 expect(workerNode.info).toStrictEqual({
549 id: expect.any(Number),
550 type: WorkerTypes.thread,
551 dynamic: false,
552 ready: true
553 })
554 }
555 })
556
557 it('Verify that pool worker tasks usage are computed', async () => {
558 const pool = new FixedClusterPool(
559 numberOfWorkers,
560 './tests/worker-files/cluster/testWorker.js'
561 )
562 const promises = new Set()
563 const maxMultiplier = 2
564 for (let i = 0; i < numberOfWorkers * maxMultiplier; i++) {
565 promises.add(pool.execute())
566 }
567 for (const workerNode of pool.workerNodes) {
568 expect(workerNode.usage).toStrictEqual({
569 tasks: {
570 executed: 0,
571 executing: maxMultiplier,
572 queued: 0,
573 maxQueued: 0,
574 failed: 0
575 },
576 runTime: {
577 history: expect.any(CircularArray)
578 },
579 waitTime: {
580 history: expect.any(CircularArray)
581 },
582 elu: {
583 idle: {
584 history: expect.any(CircularArray)
585 },
586 active: {
587 history: expect.any(CircularArray)
588 }
589 }
590 })
591 }
592 await Promise.all(promises)
593 for (const workerNode of pool.workerNodes) {
594 expect(workerNode.usage).toStrictEqual({
595 tasks: {
596 executed: maxMultiplier,
597 executing: 0,
598 queued: 0,
599 maxQueued: 0,
600 failed: 0
601 },
602 runTime: {
603 history: expect.any(CircularArray)
604 },
605 waitTime: {
606 history: expect.any(CircularArray)
607 },
608 elu: {
609 idle: {
610 history: expect.any(CircularArray)
611 },
612 active: {
613 history: expect.any(CircularArray)
614 }
615 }
616 })
617 }
618 await pool.destroy()
619 })
620
621 it('Verify that pool worker tasks usage are reset at worker choice strategy change', async () => {
622 const pool = new DynamicThreadPool(
623 Math.floor(numberOfWorkers / 2),
624 numberOfWorkers,
625 './tests/worker-files/thread/testWorker.js'
626 )
627 const promises = new Set()
628 const maxMultiplier = 2
629 for (let i = 0; i < numberOfWorkers * maxMultiplier; i++) {
630 promises.add(pool.execute())
631 }
632 await Promise.all(promises)
633 for (const workerNode of pool.workerNodes) {
634 expect(workerNode.usage).toStrictEqual({
635 tasks: {
636 executed: expect.any(Number),
637 executing: 0,
638 queued: 0,
639 maxQueued: 0,
640 failed: 0
641 },
642 runTime: {
643 history: expect.any(CircularArray)
644 },
645 waitTime: {
646 history: expect.any(CircularArray)
647 },
648 elu: {
649 idle: {
650 history: expect.any(CircularArray)
651 },
652 active: {
653 history: expect.any(CircularArray)
654 }
655 }
656 })
657 expect(workerNode.usage.tasks.executed).toBeGreaterThan(0)
658 expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(
659 numberOfWorkers * maxMultiplier
660 )
661 expect(workerNode.usage.runTime.history.length).toBe(0)
662 expect(workerNode.usage.waitTime.history.length).toBe(0)
663 expect(workerNode.usage.elu.idle.history.length).toBe(0)
664 expect(workerNode.usage.elu.active.history.length).toBe(0)
665 }
666 pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.FAIR_SHARE)
667 for (const workerNode of pool.workerNodes) {
668 expect(workerNode.usage).toStrictEqual({
669 tasks: {
670 executed: 0,
671 executing: 0,
672 queued: 0,
673 maxQueued: 0,
674 failed: 0
675 },
676 runTime: {
677 history: expect.any(CircularArray)
678 },
679 waitTime: {
680 history: expect.any(CircularArray)
681 },
682 elu: {
683 idle: {
684 history: expect.any(CircularArray)
685 },
686 active: {
687 history: expect.any(CircularArray)
688 }
689 }
690 })
691 expect(workerNode.usage.runTime.history.length).toBe(0)
692 expect(workerNode.usage.waitTime.history.length).toBe(0)
693 expect(workerNode.usage.elu.idle.history.length).toBe(0)
694 expect(workerNode.usage.elu.active.history.length).toBe(0)
695 }
696 await pool.destroy()
697 })
698
699 it("Verify that pool event emitter 'full' event can register a callback", async () => {
700 const pool = new DynamicThreadPool(
701 Math.floor(numberOfWorkers / 2),
702 numberOfWorkers,
703 './tests/worker-files/thread/testWorker.js'
704 )
705 const promises = new Set()
706 let poolFull = 0
707 let poolInfo
708 pool.emitter.on(PoolEvents.full, info => {
709 ++poolFull
710 poolInfo = info
711 })
712 for (let i = 0; i < numberOfWorkers * 2; i++) {
713 promises.add(pool.execute())
714 }
715 await Promise.all(promises)
716 // The `full` event is triggered when the number of submitted tasks at once reach the maximum number of workers in the dynamic pool.
717 // 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.
718 expect(poolFull).toBe(numberOfWorkers * 2 - 1)
719 expect(poolInfo).toStrictEqual({
720 version,
721 type: PoolTypes.dynamic,
722 worker: WorkerTypes.thread,
723 ready: expect.any(Boolean),
724 strategy: WorkerChoiceStrategies.ROUND_ROBIN,
725 minSize: expect.any(Number),
726 maxSize: expect.any(Number),
727 workerNodes: expect.any(Number),
728 idleWorkerNodes: expect.any(Number),
729 busyWorkerNodes: expect.any(Number),
730 executedTasks: expect.any(Number),
731 executingTasks: expect.any(Number),
732 queuedTasks: expect.any(Number),
733 maxQueuedTasks: expect.any(Number),
734 failedTasks: expect.any(Number)
735 })
736 await pool.destroy()
737 })
738
739 it("Verify that pool event emitter 'ready' event can register a callback", async () => {
740 const pool = new DynamicClusterPool(
741 Math.floor(numberOfWorkers / 2),
742 numberOfWorkers,
743 './tests/worker-files/cluster/testWorker.js'
744 )
745 let poolInfo
746 let poolReady = 0
747 pool.emitter.on(PoolEvents.ready, info => {
748 ++poolReady
749 poolInfo = info
750 })
751 await waitPoolEvents(pool, PoolEvents.ready, 1)
752 expect(poolReady).toBe(1)
753 expect(poolInfo).toStrictEqual({
754 version,
755 type: PoolTypes.dynamic,
756 worker: WorkerTypes.cluster,
757 ready: true,
758 strategy: WorkerChoiceStrategies.ROUND_ROBIN,
759 minSize: expect.any(Number),
760 maxSize: expect.any(Number),
761 workerNodes: expect.any(Number),
762 idleWorkerNodes: expect.any(Number),
763 busyWorkerNodes: expect.any(Number),
764 executedTasks: expect.any(Number),
765 executingTasks: expect.any(Number),
766 queuedTasks: expect.any(Number),
767 maxQueuedTasks: expect.any(Number),
768 failedTasks: expect.any(Number)
769 })
770 await pool.destroy()
771 })
772
773 it("Verify that pool event emitter 'busy' event can register a callback", async () => {
774 const pool = new FixedThreadPool(
775 numberOfWorkers,
776 './tests/worker-files/thread/testWorker.js'
777 )
778 const promises = new Set()
779 let poolBusy = 0
780 let poolInfo
781 pool.emitter.on(PoolEvents.busy, info => {
782 ++poolBusy
783 poolInfo = info
784 })
785 for (let i = 0; i < numberOfWorkers * 2; i++) {
786 promises.add(pool.execute())
787 }
788 await Promise.all(promises)
789 // The `busy` event is triggered when the number of submitted tasks at once reach the number of fixed pool workers.
790 // So in total numberOfWorkers + 1 times for a loop submitting up to numberOfWorkers * 2 tasks to the fixed pool.
791 expect(poolBusy).toBe(numberOfWorkers + 1)
792 expect(poolInfo).toStrictEqual({
793 version,
794 type: PoolTypes.fixed,
795 worker: WorkerTypes.thread,
796 ready: expect.any(Boolean),
797 strategy: WorkerChoiceStrategies.ROUND_ROBIN,
798 minSize: expect.any(Number),
799 maxSize: expect.any(Number),
800 workerNodes: expect.any(Number),
801 idleWorkerNodes: expect.any(Number),
802 busyWorkerNodes: expect.any(Number),
803 executedTasks: expect.any(Number),
804 executingTasks: expect.any(Number),
805 queuedTasks: expect.any(Number),
806 maxQueuedTasks: expect.any(Number),
807 failedTasks: expect.any(Number)
808 })
809 await pool.destroy()
810 })
811
812 it('Verify that multiple tasks worker is working', async () => {
813 const pool = new DynamicClusterPool(
814 Math.floor(numberOfWorkers / 2),
815 numberOfWorkers,
816 './tests/worker-files/cluster/testMultiTasksWorker.js'
817 )
818 const data = { n: 10 }
819 const result0 = await pool.execute(data)
820 expect(result0).toStrictEqual({ ok: 1 })
821 const result1 = await pool.execute(data, 'jsonIntegerSerialization')
822 expect(result1).toStrictEqual({ ok: 1 })
823 const result2 = await pool.execute(data, 'factorial')
824 expect(result2).toBe(3628800)
825 const result3 = await pool.execute(data, 'fibonacci')
826 expect(result3).toBe(55)
827 })
828 })