DynamicThreadPool,
FixedThreadPool,
FixedClusterPool
-} = require('../../../lib/index')
+} = require('../../../lib')
+const { CircularArray } = require('../../../lib/circular-array')
describe('Selection strategies test suite', () => {
const min = 0
it('Verify that WorkerChoiceStrategies enumeration provides string values', () => {
expect(WorkerChoiceStrategies.ROUND_ROBIN).toBe('ROUND_ROBIN')
- expect(WorkerChoiceStrategies.LESS_USED).toBe('LESS_USED')
- expect(WorkerChoiceStrategies.LESS_BUSY).toBe('LESS_BUSY')
+ expect(WorkerChoiceStrategies.LEAST_USED).toBe('LEAST_USED')
+ expect(WorkerChoiceStrategies.LEAST_BUSY).toBe('LEAST_BUSY')
+ expect(WorkerChoiceStrategies.LEAST_ELU).toBe('LEAST_ELU')
expect(WorkerChoiceStrategies.FAIR_SHARE).toBe('FAIR_SHARE')
expect(WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN).toBe(
'WEIGHTED_ROUND_ROBIN'
)
+ expect(WorkerChoiceStrategies.INTERLEAVED_WEIGHTED_ROUND_ROBIN).toBe(
+ 'INTERLEAVED_WEIGHTED_ROUND_ROBIN'
+ )
})
it('Verify ROUND_ROBIN strategy is the default at pool creation', async () => {
await pool.destroy()
})
- it('Verify ROUND_ROBIN strategy is taken at pool creation', async () => {
+ it('Verify available strategies are taken at pool creation', async () => {
+ for (const workerChoiceStrategy of Object.values(WorkerChoiceStrategies)) {
+ const pool = new FixedThreadPool(
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
+ )
+ expect(pool.opts.workerChoiceStrategy).toBe(workerChoiceStrategy)
+ expect(pool.workerChoiceStrategyContext.workerChoiceStrategy).toBe(
+ workerChoiceStrategy
+ )
+ await pool.destroy()
+ }
+ })
+
+ it('Verify available strategies can be set after pool creation', async () => {
+ for (const workerChoiceStrategy of Object.values(WorkerChoiceStrategies)) {
+ const pool = new DynamicThreadPool(
+ min,
+ max,
+ './tests/worker-files/thread/testWorker.js'
+ )
+ pool.setWorkerChoiceStrategy(workerChoiceStrategy)
+ expect(pool.opts.workerChoiceStrategy).toBe(workerChoiceStrategy)
+ expect(pool.workerChoiceStrategyContext.workerChoiceStrategy).toBe(
+ workerChoiceStrategy
+ )
+ await pool.destroy()
+ }
+ })
+
+ it('Verify available strategies default internals at pool creation', async () => {
const pool = new FixedThreadPool(
max,
- './tests/worker-files/thread/testWorker.js',
- { workerChoiceStrategy: WorkerChoiceStrategies.ROUND_ROBIN }
- )
- expect(pool.opts.workerChoiceStrategy).toBe(
- WorkerChoiceStrategies.ROUND_ROBIN
+ './tests/worker-files/thread/testWorker.js'
)
- expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.nextWorkerId
- ).toBe(0)
- // We need to clean up the resources after our test
+ for (const workerChoiceStrategy of Object.values(WorkerChoiceStrategies)) {
+ if (workerChoiceStrategy === WorkerChoiceStrategies.ROUND_ROBIN) {
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).nextWorkerNodeId
+ ).toBe(0)
+ } else if (workerChoiceStrategy === WorkerChoiceStrategies.FAIR_SHARE) {
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp
+ ).toBeInstanceOf(Array)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp.length
+ ).toBe(0)
+ } else if (
+ workerChoiceStrategy === WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN
+ ) {
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).nextWorkerNodeId
+ ).toBe(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).defaultWorkerWeight
+ ).toBeGreaterThan(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workerVirtualTaskRunTime
+ ).toBe(0)
+ }
+ }
await pool.destroy()
})
- it('Verify ROUND_ROBIN strategy can be set after pool creation', async () => {
- const pool = new DynamicThreadPool(
- min,
+ it('Verify ROUND_ROBIN strategy default policy', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.ROUND_ROBIN
+ let pool = new FixedThreadPool(
max,
- './tests/worker-files/thread/testWorker.js'
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.ROUND_ROBIN)
- expect(pool.opts.workerChoiceStrategy).toBe(
- WorkerChoiceStrategies.ROUND_ROBIN
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: true
+ })
+ await pool.destroy()
+ pool = new DynamicThreadPool(
+ min,
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: true
+ })
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify ROUND_ROBIN strategy default tasks usage statistics requirements', async () => {
+ it('Verify ROUND_ROBIN strategy default tasks statistics requirements', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.ROUND_ROBIN
let pool = new FixedThreadPool(
max,
- './tests/worker-files/thread/testWorker.js'
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.ROUND_ROBIN)
expect(
- pool.workerChoiceStrategyContext.getRequiredStatistics().runTime
- ).toBe(false)
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ waitTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: false,
+ average: false,
+ median: false
+ }
+ })
await pool.destroy()
pool = new DynamicThreadPool(
min,
max,
- './tests/worker-files/thread/testWorker.js'
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.ROUND_ROBIN)
expect(
- pool.workerChoiceStrategyContext.getRequiredStatistics().runTime
- ).toBe(false)
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ waitTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: false,
+ average: false,
+ median: false
+ }
+ })
// We need to clean up the resources after our test
await pool.destroy()
})
'./tests/worker-files/thread/testWorker.js',
{ workerChoiceStrategy: WorkerChoiceStrategies.ROUND_ROBIN }
)
- expect(pool.opts.workerChoiceStrategy).toBe(
- WorkerChoiceStrategies.ROUND_ROBIN
- )
// TODO: Create a better test to cover `RoundRobinWorkerChoiceStrategy#choose`
- const promises = []
- for (let i = 0; i < max * 2; i++) {
- promises.push(pool.execute())
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
}
await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: maxMultiplier,
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: 0
+ }
+ })
+ }
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ WorkerChoiceStrategies.ROUND_ROBIN
+ ).nextWorkerNodeId
+ ).toBe(0)
// We need to clean up the resources after our test
await pool.destroy()
})
'./tests/worker-files/thread/testWorker.js',
{ workerChoiceStrategy: WorkerChoiceStrategies.ROUND_ROBIN }
)
- expect(pool.opts.workerChoiceStrategy).toBe(
- WorkerChoiceStrategies.ROUND_ROBIN
- )
// TODO: Create a better test to cover `RoundRobinWorkerChoiceStrategy#choose`
- const promises = []
- for (let i = 0; i < max * 2; i++) {
- promises.push(pool.execute())
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
}
await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: maxMultiplier,
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: 0
+ }
+ })
+ }
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ WorkerChoiceStrategies.ROUND_ROBIN
+ ).nextWorkerNodeId
+ ).toBe(0)
// We need to clean up the resources after our test
await pool.destroy()
})
it('Verify ROUND_ROBIN strategy runtime behavior', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.ROUND_ROBIN
let pool = new FixedClusterPool(
max,
- './tests/worker-files/cluster/testWorker.js'
+ './tests/worker-files/cluster/testWorker.js',
+ { workerChoiceStrategy }
)
let results = new Set()
for (let i = 0; i < max; i++) {
- results.add(pool.chooseWorker()[1].id)
+ results.add(pool.workerNodes[pool.chooseWorkerNode()].worker.id)
}
expect(results.size).toBe(max)
await pool.destroy()
- pool = new FixedThreadPool(max, './tests/worker-files/thread/testWorker.js')
+ pool = new FixedThreadPool(
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
+ )
results = new Set()
for (let i = 0; i < max; i++) {
- results.add(pool.chooseWorker()[1].threadId)
+ results.add(pool.workerNodes[pool.chooseWorkerNode()].worker.threadId)
}
expect(results.size).toBe(max)
await pool.destroy()
})
it('Verify ROUND_ROBIN strategy internals are resets after setting it', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.ROUND_ROBIN
let pool = new FixedThreadPool(
max,
'./tests/worker-files/thread/testWorker.js',
{ workerChoiceStrategy: WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN }
)
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.nextWorkerId
- ).toBeUndefined()
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.ROUND_ROBIN)
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).nextWorkerNodeId
+ ).toBeDefined()
+ pool.setWorkerChoiceStrategy(workerChoiceStrategy)
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.nextWorkerId
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).nextWorkerNodeId
).toBe(0)
await pool.destroy()
pool = new DynamicThreadPool(
{ workerChoiceStrategy: WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN }
)
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy
- .nextWorkerId
- ).toBeUndefined()
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.ROUND_ROBIN)
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).nextWorkerNodeId
+ ).toBeDefined()
+ pool.setWorkerChoiceStrategy(workerChoiceStrategy)
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy
- .nextWorkerId
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).nextWorkerNodeId
).toBe(0)
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify LESS_USED strategy is taken at pool creation', async () => {
- const pool = new FixedThreadPool(
+ it('Verify LEAST_USED strategy default policy', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.LEAST_USED
+ let pool = new FixedThreadPool(
max,
'./tests/worker-files/thread/testWorker.js',
- { workerChoiceStrategy: WorkerChoiceStrategies.LESS_USED }
- )
- expect(pool.opts.workerChoiceStrategy).toBe(
- WorkerChoiceStrategies.LESS_USED
+ { workerChoiceStrategy }
)
- // We need to clean up the resources after our test
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: false
+ })
await pool.destroy()
- })
-
- it('Verify LESS_USED strategy can be set after pool creation', async () => {
- const pool = new FixedThreadPool(
+ pool = new DynamicThreadPool(
+ min,
max,
- './tests/worker-files/thread/testWorker.js'
- )
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.LESS_USED)
- expect(pool.opts.workerChoiceStrategy).toBe(
- WorkerChoiceStrategies.LESS_USED
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: false
+ })
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify LESS_USED strategy default tasks usage statistics requirements', async () => {
+ it('Verify LEAST_USED strategy default tasks statistics requirements', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.LEAST_USED
let pool = new FixedThreadPool(
max,
- './tests/worker-files/thread/testWorker.js'
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.LESS_USED)
expect(
- pool.workerChoiceStrategyContext.getRequiredStatistics().runTime
- ).toBe(false)
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ waitTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: false,
+ average: false,
+ median: false
+ }
+ })
await pool.destroy()
pool = new DynamicThreadPool(
min,
max,
- './tests/worker-files/thread/testWorker.js'
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.LESS_USED)
expect(
- pool.workerChoiceStrategyContext.getRequiredStatistics().runTime
- ).toBe(false)
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ waitTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: false,
+ average: false,
+ median: false
+ }
+ })
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify LESS_USED strategy can be run in a fixed pool', async () => {
+ it('Verify LEAST_USED strategy can be run in a fixed pool', async () => {
const pool = new FixedThreadPool(
max,
'./tests/worker-files/thread/testWorker.js',
- { workerChoiceStrategy: WorkerChoiceStrategies.LESS_USED }
+ { workerChoiceStrategy: WorkerChoiceStrategies.LEAST_USED }
)
- // TODO: Create a better test to cover `LessUsedWorkerChoiceStrategy#choose`
- const promises = []
- for (let i = 0; i < max * 2; i++) {
- promises.push(pool.execute())
+ // TODO: Create a better test to cover `LeastUsedWorkerChoiceStrategy#choose`
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
}
await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: expect.any(Number),
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: 0
+ }
+ })
+ expect(workerNode.usage.tasks.executed).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(
+ max * maxMultiplier
+ )
+ }
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify LESS_USED strategy can be run in a dynamic pool', async () => {
+ it('Verify LEAST_USED strategy can be run in a dynamic pool', async () => {
const pool = new DynamicThreadPool(
min,
max,
'./tests/worker-files/thread/testWorker.js',
- { workerChoiceStrategy: WorkerChoiceStrategies.LESS_USED }
+ { workerChoiceStrategy: WorkerChoiceStrategies.LEAST_USED }
)
- // TODO: Create a better test to cover `LessUsedWorkerChoiceStrategy#choose`
- const promises = []
- for (let i = 0; i < max * 2; i++) {
- promises.push(pool.execute())
+ // TODO: Create a better test to cover `LeastUsedWorkerChoiceStrategy#choose`
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
}
await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: expect.any(Number),
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: 0
+ }
+ })
+ expect(workerNode.usage.tasks.executed).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(
+ max * maxMultiplier
+ )
+ }
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify LESS_BUSY strategy is taken at pool creation', async () => {
- const pool = new FixedThreadPool(
+ it('Verify LEAST_BUSY strategy default policy', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.LEAST_BUSY
+ let pool = new FixedThreadPool(
max,
'./tests/worker-files/thread/testWorker.js',
- { workerChoiceStrategy: WorkerChoiceStrategies.LESS_BUSY }
- )
- expect(pool.opts.workerChoiceStrategy).toBe(
- WorkerChoiceStrategies.LESS_BUSY
+ { workerChoiceStrategy }
)
- // We need to clean up the resources after our test
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: false
+ })
await pool.destroy()
- })
-
- it('Verify LESS_BUSY strategy can be set after pool creation', async () => {
- const pool = new FixedThreadPool(
+ pool = new DynamicThreadPool(
+ min,
max,
- './tests/worker-files/thread/testWorker.js'
- )
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.LESS_BUSY)
- expect(pool.opts.workerChoiceStrategy).toBe(
- WorkerChoiceStrategies.LESS_BUSY
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: false
+ })
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify LESS_BUSY strategy default tasks usage statistics requirements', async () => {
+ it('Verify LEAST_BUSY strategy default tasks statistics requirements', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.LEAST_BUSY
let pool = new FixedThreadPool(
max,
- './tests/worker-files/thread/testWorker.js'
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.LESS_BUSY)
expect(
- pool.workerChoiceStrategyContext.getRequiredStatistics().runTime
- ).toBe(true)
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: true,
+ average: false,
+ median: false
+ },
+ waitTime: {
+ aggregate: true,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: false,
+ average: false,
+ median: false
+ }
+ })
await pool.destroy()
pool = new DynamicThreadPool(
min,
max,
- './tests/worker-files/thread/testWorker.js'
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.LESS_BUSY)
expect(
- pool.workerChoiceStrategyContext.getRequiredStatistics().runTime
- ).toBe(true)
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: true,
+ average: false,
+ median: false
+ },
+ waitTime: {
+ aggregate: true,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: false,
+ average: false,
+ median: false
+ }
+ })
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify LESS_BUSY strategy can be run in a fixed pool', async () => {
+ it('Verify LEAST_BUSY strategy can be run in a fixed pool', async () => {
const pool = new FixedThreadPool(
max,
'./tests/worker-files/thread/testWorker.js',
- { workerChoiceStrategy: WorkerChoiceStrategies.LESS_BUSY }
+ { workerChoiceStrategy: WorkerChoiceStrategies.LEAST_BUSY }
)
- // TODO: Create a better test to cover `LessBusyWorkerChoiceStrategy#choose`
- const promises = []
- for (let i = 0; i < max * 2; i++) {
- promises.push(pool.execute())
+ // TODO: Create a better test to cover `LeastBusyWorkerChoiceStrategy#choose`
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
}
await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: expect.any(Number),
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: expect.any(Number),
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: expect.any(Number),
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: 0
+ }
+ })
+ expect(workerNode.usage.tasks.executed).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(
+ max * maxMultiplier
+ )
+ expect(workerNode.usage.runTime.aggregate).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.waitTime.aggregate).toBeGreaterThanOrEqual(0)
+ }
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify LESS_BUSY strategy can be run in a dynamic pool', async () => {
+ it('Verify LEAST_BUSY strategy can be run in a dynamic pool', async () => {
const pool = new DynamicThreadPool(
min,
max,
'./tests/worker-files/thread/testWorker.js',
- { workerChoiceStrategy: WorkerChoiceStrategies.LESS_BUSY }
+ { workerChoiceStrategy: WorkerChoiceStrategies.LEAST_BUSY }
)
- // TODO: Create a better test to cover `LessBusyWorkerChoiceStrategy#choose`
- const promises = []
- for (let i = 0; i < max * 2; i++) {
- promises.push(pool.execute())
+ // TODO: Create a better test to cover `LeastBusyWorkerChoiceStrategy#choose`
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
}
await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: expect.any(Number),
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: expect.any(Number),
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: expect.any(Number),
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: 0
+ }
+ })
+ expect(workerNode.usage.tasks.executed).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(
+ max * maxMultiplier
+ )
+ expect(workerNode.usage.runTime.aggregate).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.waitTime.aggregate).toBeGreaterThanOrEqual(0)
+ }
+ // We need to clean up the resources after our test
+ await pool.destroy()
+ })
+
+ it('Verify LEAST_ELU strategy default policy', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.LEAST_ELU
+ let pool = new FixedThreadPool(
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
+ )
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: false
+ })
+ await pool.destroy()
+ pool = new DynamicThreadPool(
+ min,
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
+ )
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: false
+ })
+ // We need to clean up the resources after our test
+ await pool.destroy()
+ })
+
+ it('Verify LEAST_ELU strategy default tasks statistics requirements', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.LEAST_ELU
+ let pool = new FixedThreadPool(
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
+ )
+ expect(
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ waitTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: true,
+ average: false,
+ median: false
+ }
+ })
+ await pool.destroy()
+ pool = new DynamicThreadPool(
+ min,
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
+ )
+ expect(
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ waitTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: true,
+ average: false,
+ median: false
+ }
+ })
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify FAIR_SHARE strategy is taken at pool creation', async () => {
+ it('Verify LEAST_ELU strategy can be run in a fixed pool', async () => {
const pool = new FixedThreadPool(
max,
'./tests/worker-files/thread/testWorker.js',
- { workerChoiceStrategy: WorkerChoiceStrategies.FAIR_SHARE }
+ { workerChoiceStrategy: WorkerChoiceStrategies.LEAST_ELU }
)
- expect(pool.opts.workerChoiceStrategy).toBe(
- WorkerChoiceStrategies.FAIR_SHARE
+ // TODO: Create a better test to cover `LeastEluWorkerChoiceStrategy#choose`
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
+ }
+ await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: expect.any(Number),
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: expect.any(Number),
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: expect.any(Number)
+ }
+ })
+ expect(workerNode.usage.tasks.executed).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(
+ max * maxMultiplier
+ )
+ expect(workerNode.usage.elu.utilization).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.elu.utilization).toBeLessThanOrEqual(1)
+ }
+ // We need to clean up the resources after our test
+ await pool.destroy()
+ })
+
+ it('Verify LEAST_ELU strategy can be run in a dynamic pool', async () => {
+ const pool = new DynamicThreadPool(
+ min,
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy: WorkerChoiceStrategies.LEAST_ELU }
)
- for (const workerKey of pool.workerChoiceStrategyContext.workerChoiceStrategy.workerLastVirtualTaskTimestamp.keys()) {
- expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerLastVirtualTaskTimestamp.get(
- workerKey
- ).start
- ).toBe(0)
- expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerLastVirtualTaskTimestamp.get(
- workerKey
- ).end
- ).toBe(0)
+ // TODO: Create a better test to cover `LeastEluWorkerChoiceStrategy#choose`
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
+ }
+ await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: expect.any(Number),
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: expect.any(Number),
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: expect.any(Number)
+ }
+ })
+ expect(workerNode.usage.tasks.executed).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(
+ max * maxMultiplier
+ )
+ expect(workerNode.usage.elu.utilization).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.elu.utilization).toBeLessThanOrEqual(1)
}
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify FAIR_SHARE strategy can be set after pool creation', async () => {
- const pool = new FixedThreadPool(
+ it('Verify FAIR_SHARE strategy default policy', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.FAIR_SHARE
+ let pool = new FixedThreadPool(
max,
- './tests/worker-files/thread/testWorker.js'
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.FAIR_SHARE)
- expect(pool.opts.workerChoiceStrategy).toBe(
- WorkerChoiceStrategies.FAIR_SHARE
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: false
+ })
+ await pool.destroy()
+ pool = new DynamicThreadPool(
+ min,
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: false
+ })
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify FAIR_SHARE strategy default tasks usage statistics requirements', async () => {
+ it('Verify FAIR_SHARE strategy default tasks statistics requirements', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.FAIR_SHARE
let pool = new FixedThreadPool(
max,
- './tests/worker-files/thread/testWorker.js'
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.FAIR_SHARE)
expect(
- pool.workerChoiceStrategyContext.getRequiredStatistics().runTime
- ).toBe(true)
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: true,
+ average: true,
+ median: false
+ },
+ waitTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: true,
+ average: true,
+ median: false
+ }
+ })
await pool.destroy()
pool = new DynamicThreadPool(
min,
max,
- './tests/worker-files/thread/testWorker.js'
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.FAIR_SHARE)
expect(
- pool.workerChoiceStrategyContext.getRequiredStatistics().runTime
- ).toBe(true)
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: true,
+ average: true,
+ median: false
+ },
+ waitTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: true,
+ average: true,
+ median: false
+ }
+ })
// We need to clean up the resources after our test
await pool.destroy()
})
{ workerChoiceStrategy: WorkerChoiceStrategies.FAIR_SHARE }
)
// TODO: Create a better test to cover `FairShareChoiceStrategy#choose`
- const promises = []
- for (let i = 0; i < max * 2; i++) {
- promises.push(pool.execute())
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
}
await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: expect.any(Number),
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: expect.any(Number),
+ average: expect.any(Number),
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: expect.any(Number),
+ average: expect.any(Number),
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: expect.any(Number)
+ }
+ })
+ expect(workerNode.usage.tasks.executed).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(
+ max * maxMultiplier
+ )
+ expect(workerNode.usage.runTime.aggregate).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.runTime.average).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.elu.utilization).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.elu.utilization).toBeLessThanOrEqual(1)
+ }
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy
- .workerLastVirtualTaskTimestamp.size
- ).toBe(pool.workers.length)
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp.length
+ ).toBe(pool.workerNodes.length)
// We need to clean up the resources after our test
await pool.destroy()
})
{ workerChoiceStrategy: WorkerChoiceStrategies.FAIR_SHARE }
)
// TODO: Create a better test to cover `FairShareChoiceStrategy#choose`
- const promises = []
+ const promises = new Set()
const maxMultiplier = 2
for (let i = 0; i < max * maxMultiplier; i++) {
- promises.push(pool.execute())
+ promises.add(pool.execute())
}
await Promise.all(promises)
- // if (process.platform !== 'win32') {
- // expect(
- // pool.workerChoiceStrategyContext.workerChoiceStrategy
- // .workerChoiceStrategy.workerLastVirtualTaskTimestamp.size
- // ).toBe(pool.workers.length)
- // }
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: expect.any(Number),
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: expect.any(Number),
+ average: expect.any(Number),
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: expect.any(Number),
+ average: expect.any(Number),
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: expect.any(Number)
+ }
+ })
+ expect(workerNode.usage.tasks.executed).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(
+ max * maxMultiplier
+ )
+ expect(workerNode.usage.runTime.aggregate).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.runTime.average).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.elu.utilization).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.elu.utilization).toBeLessThanOrEqual(1)
+ }
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp.length
+ ).toBe(pool.workerNodes.length)
+ // We need to clean up the resources after our test
+ await pool.destroy()
+ })
+
+ it('Verify FAIR_SHARE strategy can be run in a dynamic pool with median runtime statistic', async () => {
+ const pool = new DynamicThreadPool(
+ min,
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ {
+ workerChoiceStrategy: WorkerChoiceStrategies.FAIR_SHARE,
+ workerChoiceStrategyOptions: {
+ runTime: { median: true }
+ }
+ }
+ )
+ // TODO: Create a better test to cover `FairShareChoiceStrategy#choose`
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
+ }
+ await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: expect.any(Number),
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: expect.any(Number),
+ average: 0,
+ median: expect.any(Number),
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: expect.any(Number),
+ average: expect.any(Number),
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: expect.any(Number)
+ }
+ })
+ expect(workerNode.usage.tasks.executed).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(
+ max * maxMultiplier
+ )
+ expect(workerNode.usage.runTime.aggregate).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.runTime.median).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.elu.utilization).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.elu.utilization).toBeLessThanOrEqual(1)
+ }
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp.length
+ ).toBe(pool.workerNodes.length)
// We need to clean up the resources after our test
await pool.destroy()
})
it('Verify FAIR_SHARE strategy internals are resets after setting it', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.FAIR_SHARE
let pool = new FixedThreadPool(
max,
'./tests/worker-files/thread/testWorker.js'
)
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy
- .workerLastVirtualTaskTimestamp
- ).toBeUndefined()
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.FAIR_SHARE)
- for (const workerKey of pool.workerChoiceStrategyContext.workerChoiceStrategy.workerLastVirtualTaskTimestamp.keys()) {
- expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerLastVirtualTaskTimestamp.get(
- workerKey
- ).start
- ).toBe(0)
- expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerLastVirtualTaskTimestamp.get(
- workerKey
- ).end
- ).toBe(0)
- }
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp
+ ).toBeInstanceOf(Array)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp.length
+ ).toBe(0)
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp[0] = performance.now()
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp.length
+ ).toBe(1)
+ pool.setWorkerChoiceStrategy(workerChoiceStrategy)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp
+ ).toBeInstanceOf(Array)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp.length
+ ).toBe(0)
await pool.destroy()
pool = new DynamicThreadPool(
min,
'./tests/worker-files/thread/testWorker.js'
)
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy
- .workerLastVirtualTaskTimestamp
- ).toBeUndefined()
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.FAIR_SHARE)
- for (const workerKey of pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy.workerLastVirtualTaskTimestamp.keys()) {
- expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy.workerLastVirtualTaskTimestamp.get(
- workerKey
- ).start
- ).toBe(0)
- expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy.workerLastVirtualTaskTimestamp.get(
- workerKey
- ).end
- ).toBe(0)
- }
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp
+ ).toBeInstanceOf(Array)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp.length
+ ).toBe(0)
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp[0] = performance.now()
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp.length
+ ).toBe(1)
+ pool.setWorkerChoiceStrategy(workerChoiceStrategy)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp
+ ).toBeInstanceOf(Array)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workersVirtualTaskEndTimestamp.length
+ ).toBe(0)
+ // We need to clean up the resources after our test
+ await pool.destroy()
+ })
+
+ it('Verify WEIGHTED_ROUND_ROBIN strategy default policy', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN
+ let pool = new FixedThreadPool(
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
+ )
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: true
+ })
+ await pool.destroy()
+ pool = new DynamicThreadPool(
+ min,
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
+ )
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: true
+ })
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify WEIGHTED_ROUND_ROBIN strategy is taken at pool creation', async () => {
+ it('Verify WEIGHTED_ROUND_ROBIN strategy default tasks statistics requirements', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN
+ let pool = new FixedThreadPool(
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
+ )
+ expect(
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: true,
+ average: true,
+ median: false
+ },
+ waitTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: false,
+ average: false,
+ median: false
+ }
+ })
+ await pool.destroy()
+ pool = new DynamicThreadPool(
+ min,
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
+ )
+ expect(
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: true,
+ average: true,
+ median: false
+ },
+ waitTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: false,
+ average: false,
+ median: false
+ }
+ })
+ // We need to clean up the resources after our test
+ await pool.destroy()
+ })
+
+ it('Verify WEIGHTED_ROUND_ROBIN strategy can be run in a fixed pool', async () => {
const pool = new FixedThreadPool(
max,
'./tests/worker-files/thread/testWorker.js',
{ workerChoiceStrategy: WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN }
)
- expect(pool.opts.workerChoiceStrategy).toBe(
- WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN
- )
+ // TODO: Create a better test to cover `WeightedRoundRobinWorkerChoiceStrategy#choose`
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
+ }
+ await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: expect.any(Number),
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: expect.any(Number),
+ average: expect.any(Number),
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: 0
+ }
+ })
+ expect(workerNode.usage.tasks.executed).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(
+ max * maxMultiplier
+ )
+ expect(workerNode.usage.runTime.aggregate).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.runTime.average).toBeGreaterThanOrEqual(0)
+ }
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.currentWorkerId
- ).toBe(0)
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).defaultWorkerWeight
+ ).toBeGreaterThan(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).workerVirtualTaskRunTime
+ ).toBeGreaterThanOrEqual(0)
+ // We need to clean up the resources after our test
+ await pool.destroy()
+ })
+
+ it('Verify WEIGHTED_ROUND_ROBIN strategy can be run in a dynamic pool', async () => {
+ const pool = new DynamicThreadPool(
+ min,
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy: WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN }
+ )
+ // TODO: Create a better test to cover `WeightedRoundRobinWorkerChoiceStrategy#choose`
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
+ }
+ await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: expect.any(Number),
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: expect.any(Number),
+ average: expect.any(Number),
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: 0
+ }
+ })
+ expect(workerNode.usage.tasks.executed).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(
+ max * maxMultiplier
+ )
+ expect(workerNode.usage.runTime.aggregate).toBeGreaterThan(0)
+ expect(workerNode.usage.runTime.average).toBeGreaterThan(0)
+ }
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.defaultWorkerWeight
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).defaultWorkerWeight
).toBeGreaterThan(0)
- for (const workerKey of pool.workerChoiceStrategyContext.workerChoiceStrategy.workersTaskRunTime.keys()) {
- expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workersTaskRunTime.get(
- workerKey
- ).weight
- ).toBeGreaterThan(0)
- expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workersTaskRunTime.get(
- workerKey
- ).runTime
- ).toBe(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).workerVirtualTaskRunTime
+ ).toBeGreaterThanOrEqual(0)
+ // We need to clean up the resources after our test
+ await pool.destroy()
+ })
+
+ it('Verify WEIGHTED_ROUND_ROBIN strategy can be run in a dynamic pool with median runtime statistic', async () => {
+ const pool = new DynamicThreadPool(
+ min,
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ {
+ workerChoiceStrategy: WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN,
+ workerChoiceStrategyOptions: {
+ runTime: { median: true }
+ }
+ }
+ )
+ // TODO: Create a better test to cover `WeightedRoundRobinWorkerChoiceStrategy#choose`
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
+ }
+ await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: expect.any(Number),
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: expect.any(Number),
+ average: 0,
+ median: expect.any(Number),
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: 0
+ }
+ })
+ expect(workerNode.usage.tasks.executed).toBeGreaterThanOrEqual(0)
+ expect(workerNode.usage.tasks.executed).toBeLessThanOrEqual(
+ max * maxMultiplier
+ )
+ expect(workerNode.usage.runTime.aggregate).toBeGreaterThan(0)
+ expect(workerNode.usage.runTime.median).toBeGreaterThan(0)
}
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).defaultWorkerWeight
+ ).toBeGreaterThan(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).workerVirtualTaskRunTime
+ ).toBeGreaterThanOrEqual(0)
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify WEIGHTED_ROUND_ROBIN strategy can be set after pool creation', async () => {
- const pool = new FixedThreadPool(
+ it('Verify WEIGHTED_ROUND_ROBIN strategy internals are resets after setting it', async () => {
+ const workerChoiceStrategy = WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN
+ let pool = new FixedThreadPool(
max,
'./tests/worker-files/thread/testWorker.js'
)
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN)
- expect(pool.opts.workerChoiceStrategy).toBe(
- WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).nextWorkerNodeId
+ ).toBeDefined()
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).defaultWorkerWeight
+ ).toBeDefined()
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workerVirtualTaskRunTime
+ ).toBeDefined()
+ pool.setWorkerChoiceStrategy(workerChoiceStrategy)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).nextWorkerNodeId
+ ).toBe(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).defaultWorkerWeight
+ ).toBeGreaterThan(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workerVirtualTaskRunTime
+ ).toBe(0)
+ await pool.destroy()
+ pool = new DynamicThreadPool(
+ min,
+ max,
+ './tests/worker-files/thread/testWorker.js'
+ )
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).nextWorkerNodeId
+ ).toBeDefined()
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).defaultWorkerWeight
+ ).toBeDefined()
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workerVirtualTaskRunTime
+ ).toBeDefined()
+ pool.setWorkerChoiceStrategy(workerChoiceStrategy)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).nextWorkerNodeId
+ ).toBe(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).defaultWorkerWeight
+ ).toBeGreaterThan(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).workerVirtualTaskRunTime
+ ).toBe(0)
+ // We need to clean up the resources after our test
+ await pool.destroy()
+ })
+
+ it('Verify INTERLEAVED_WEIGHTED_ROUND_ROBIN strategy default policy', async () => {
+ const workerChoiceStrategy =
+ WorkerChoiceStrategies.INTERLEAVED_WEIGHTED_ROUND_ROBIN
+ let pool = new FixedThreadPool(
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: true
+ })
+ await pool.destroy()
+ pool = new DynamicThreadPool(
+ min,
+ max,
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
+ )
+ expect(pool.workerChoiceStrategyContext.getStrategyPolicy()).toStrictEqual({
+ useDynamicWorker: true
+ })
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify WEIGHTED_ROUND_ROBIN strategy default tasks usage statistics requirements', async () => {
+ it('Verify INTERLEAVED_WEIGHTED_ROUND_ROBIN strategy default tasks statistics requirements', async () => {
+ const workerChoiceStrategy =
+ WorkerChoiceStrategies.INTERLEAVED_WEIGHTED_ROUND_ROBIN
let pool = new FixedThreadPool(
max,
- './tests/worker-files/thread/testWorker.js'
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN)
expect(
- pool.workerChoiceStrategyContext.getRequiredStatistics().runTime
- ).toBe(true)
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ waitTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: false,
+ average: false,
+ median: false
+ }
+ })
await pool.destroy()
pool = new DynamicThreadPool(
min,
max,
- './tests/worker-files/thread/testWorker.js'
+ './tests/worker-files/thread/testWorker.js',
+ { workerChoiceStrategy }
)
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN)
expect(
- pool.workerChoiceStrategyContext.getRequiredStatistics().runTime
- ).toBe(true)
+ pool.workerChoiceStrategyContext.getTaskStatisticsRequirements()
+ ).toStrictEqual({
+ runTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ waitTime: {
+ aggregate: false,
+ average: false,
+ median: false
+ },
+ elu: {
+ aggregate: false,
+ average: false,
+ median: false
+ }
+ })
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify WEIGHTED_ROUND_ROBIN strategy can be run in a fixed pool', async () => {
+ it('Verify INTERLEAVED_WEIGHTED_ROUND_ROBIN strategy can be run in a fixed pool', async () => {
const pool = new FixedThreadPool(
max,
'./tests/worker-files/thread/testWorker.js',
- { workerChoiceStrategy: WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN }
+ {
+ workerChoiceStrategy:
+ WorkerChoiceStrategies.INTERLEAVED_WEIGHTED_ROUND_ROBIN
+ }
)
- // TODO: Create a better test to cover `WeightedRoundRobinWorkerChoiceStrategy#choose`
- const promises = []
- for (let i = 0; i < max * 2; i++) {
- promises.push(pool.execute())
+ // TODO: Create a better test to cover `InterleavedWeightedRoundRobinWorkerChoiceStrategy#choose`
+ const promises = new Set()
+ const maxMultiplier = 2
+ for (let i = 0; i < max * maxMultiplier; i++) {
+ promises.add(pool.execute())
}
await Promise.all(promises)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: maxMultiplier,
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: 0
+ }
+ })
+ }
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).defaultWorkerWeight
+ ).toBeGreaterThan(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).roundId
+ ).toBe(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).nextWorkerNodeId
+ ).toBe(0)
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workersTaskRunTime
- .size
- ).toBe(pool.workers.length)
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).roundWeights
+ ).toStrictEqual([
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).defaultWorkerWeight
+ ])
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify WEIGHTED_ROUND_ROBIN strategy can be run in a dynamic pool', async () => {
+ it('Verify INTERLEAVED_WEIGHTED_ROUND_ROBIN strategy can be run in a dynamic pool', async () => {
const pool = new DynamicThreadPool(
min,
max,
'./tests/worker-files/thread/testWorker.js',
- { workerChoiceStrategy: WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN }
+ {
+ workerChoiceStrategy:
+ WorkerChoiceStrategies.INTERLEAVED_WEIGHTED_ROUND_ROBIN
+ }
)
- // TODO: Create a better test to cover `WeightedRoundRobinWorkerChoiceStrategy#choose`
- const promises = []
- const maxMultiplier =
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy
- .defaultWorkerWeight * 2
+ // TODO: Create a better test to cover `InterleavedWeightedRoundRobinWorkerChoiceStrategy#choose`
+ const promises = new Set()
+ const maxMultiplier = 2
for (let i = 0; i < max * maxMultiplier; i++) {
- promises.push(pool.execute())
+ promises.add(pool.execute())
}
await Promise.all(promises)
- if (process.platform !== 'win32') {
- expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy
- .workerChoiceStrategy.workersTaskRunTime.size
- ).toBe(pool.workers.length)
+ for (const workerNode of pool.workerNodes) {
+ expect(workerNode.usage).toStrictEqual({
+ tasks: {
+ executed: maxMultiplier,
+ executing: 0,
+ queued: 0,
+ maxQueued: 0,
+ failed: 0
+ },
+ runTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ waitTime: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ elu: {
+ idle: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ active: {
+ aggregate: 0,
+ average: 0,
+ median: 0,
+ history: expect.any(CircularArray)
+ },
+ utilization: 0
+ }
+ })
}
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).defaultWorkerWeight
+ ).toBeGreaterThan(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).roundId
+ ).toBe(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).nextWorkerNodeId
+ ).toBe(0)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).roundWeights
+ ).toStrictEqual([
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).defaultWorkerWeight
+ ])
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify WEIGHTED_ROUND_ROBIN strategy internals are resets after setting it', async () => {
+ it('Verify INTERLEAVED_WEIGHTED_ROUND_ROBIN strategy internals are resets after setting it', async () => {
+ const workerChoiceStrategy =
+ WorkerChoiceStrategies.INTERLEAVED_WEIGHTED_ROUND_ROBIN
let pool = new FixedThreadPool(
max,
'./tests/worker-files/thread/testWorker.js'
)
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.currentWorkerId
- ).toBeUndefined()
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).roundId
+ ).toBeDefined()
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.defaultWorkerWeight
- ).toBeUndefined()
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).nextWorkerNodeId
+ ).toBeDefined()
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workersTaskRunTime
- ).toBeUndefined()
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN)
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).defaultWorkerWeight
+ ).toBeDefined()
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).roundWeights
+ ).toBeDefined()
+ pool.setWorkerChoiceStrategy(workerChoiceStrategy)
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).roundId
+ ).toBe(0)
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.currentWorkerId
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).nextWorkerNodeId
).toBe(0)
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.defaultWorkerWeight
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).defaultWorkerWeight
).toBeGreaterThan(0)
- for (const workerKey of pool.workerChoiceStrategyContext.workerChoiceStrategy.workersTaskRunTime.keys()) {
- expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workersTaskRunTime.get(
- workerKey
- ).runTime
- ).toBe(0)
- }
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).roundWeights
+ ).toStrictEqual([
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).defaultWorkerWeight
+ ])
await pool.destroy()
pool = new DynamicThreadPool(
min,
'./tests/worker-files/thread/testWorker.js'
)
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy
- .currentWorkerId
- ).toBeUndefined()
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).roundId
+ ).toBeDefined()
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).nextWorkerNodeId
+ ).toBeDefined()
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy
- .defaultWorkerWeight
- ).toBeUndefined()
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).defaultWorkerWeight
+ ).toBeDefined()
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy
- .workersTaskRunTime
- ).toBeUndefined()
- pool.setWorkerChoiceStrategy(WorkerChoiceStrategies.WEIGHTED_ROUND_ROBIN)
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).roundWeights
+ ).toBeDefined()
+ pool.setWorkerChoiceStrategy(workerChoiceStrategy)
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy
- .currentWorkerId
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).nextWorkerNodeId
).toBe(0)
expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy
- .defaultWorkerWeight
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).defaultWorkerWeight
).toBeGreaterThan(0)
- for (const workerKey of pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy.workersTaskRunTime.keys()) {
- expect(
- pool.workerChoiceStrategyContext.workerChoiceStrategy.workerChoiceStrategy.workersTaskRunTime.get(
- workerKey
- ).runTime
- ).toBe(0)
- }
+ expect(
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ workerChoiceStrategy
+ ).roundWeights
+ ).toStrictEqual([
+ pool.workerChoiceStrategyContext.workerChoiceStrategies.get(
+ pool.workerChoiceStrategyContext.workerChoiceStrategy
+ ).defaultWorkerWeight
+ ])
// We need to clean up the resources after our test
await pool.destroy()
})
- it('Verify unknown strategies throw error', () => {
+ it('Verify unknown strategy throw error', () => {
expect(
() =>
new DynamicThreadPool(
'./tests/worker-files/thread/testWorker.js',
{ workerChoiceStrategy: 'UNKNOWN_STRATEGY' }
)
- ).toThrowError(
- new Error("Worker choice strategy 'UNKNOWN_STRATEGY' not found")
- )
+ ).toThrowError("Invalid worker choice strategy 'UNKNOWN_STRATEGY'")
})
})