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