|
| 1 | +import unittest |
| 2 | +from executorlib import SingleNodeExecutor |
| 3 | +from executorlib.standalone.serialize import cloudpickle_register |
| 4 | + |
| 5 | + |
| 6 | +def sleep_funct(sec): |
| 7 | + from time import sleep |
| 8 | + sleep(sec) |
| 9 | + return sec |
| 10 | + |
| 11 | + |
| 12 | +class TestResizing(unittest.TestCase): |
| 13 | + def test_without_dependencies_decrease(self): |
| 14 | + cloudpickle_register(ind=1) |
| 15 | + with SingleNodeExecutor(max_workers=2, block_allocation=True, disable_dependencies=True) as exe: |
| 16 | + future_lst = [exe.submit(sleep_funct, 1) for _ in range(4)] |
| 17 | + self.assertEqual([f.done() for f in future_lst], [False, False, False, False]) |
| 18 | + self.assertEqual(len(exe), 4) |
| 19 | + sleep_funct(sec=0.5) |
| 20 | + exe.max_workers = 1 |
| 21 | + self.assertTrue(len(exe) >= 1) |
| 22 | + self.assertEqual(len(exe._process), 1) |
| 23 | + self.assertTrue(1 <= sum([f.done() for f in future_lst]) < 3) |
| 24 | + self.assertEqual([f.result() for f in future_lst], [1, 1, 1, 1]) |
| 25 | + self.assertEqual([f.done() for f in future_lst], [True, True, True, True]) |
| 26 | + |
| 27 | + def test_without_dependencies_increase(self): |
| 28 | + cloudpickle_register(ind=1) |
| 29 | + with SingleNodeExecutor(max_workers=1, block_allocation=True, disable_dependencies=True) as exe: |
| 30 | + future_lst = [exe.submit(sleep_funct, 0.1) for _ in range(4)] |
| 31 | + self.assertEqual([f.done() for f in future_lst], [False, False, False, False]) |
| 32 | + self.assertEqual(len(exe), 4) |
| 33 | + self.assertEqual(exe.max_workers, 1) |
| 34 | + future_lst[0].result() |
| 35 | + exe.max_workers = 2 |
| 36 | + self.assertEqual(exe.max_workers, 2) |
| 37 | + self.assertTrue(len(exe) >= 1) |
| 38 | + self.assertEqual(len(exe._process), 2) |
| 39 | + self.assertEqual([f.done() for f in future_lst], [True, False, False, False]) |
| 40 | + self.assertEqual([f.result() for f in future_lst], [0.1, 0.1, 0.1, 0.1]) |
| 41 | + self.assertEqual([f.done() for f in future_lst], [True, True, True, True]) |
| 42 | + |
| 43 | + def test_with_dependencies_decrease(self): |
| 44 | + cloudpickle_register(ind=1) |
| 45 | + with SingleNodeExecutor(max_workers=2, block_allocation=True, disable_dependencies=False) as exe: |
| 46 | + future_lst = [exe.submit(sleep_funct, 1) for _ in range(4)] |
| 47 | + self.assertEqual([f.done() for f in future_lst], [False, False, False, False]) |
| 48 | + self.assertEqual(len(exe), 4) |
| 49 | + sleep_funct(sec=0.5) |
| 50 | + exe.max_workers = 1 |
| 51 | + self.assertTrue(1 <= sum([f.done() for f in future_lst]) < 3) |
| 52 | + self.assertEqual([f.result() for f in future_lst], [1, 1, 1, 1]) |
| 53 | + self.assertEqual([f.done() for f in future_lst], [True, True, True, True]) |
| 54 | + |
| 55 | + def test_with_dependencies_increase(self): |
| 56 | + cloudpickle_register(ind=1) |
| 57 | + with SingleNodeExecutor(max_workers=1, block_allocation=True, disable_dependencies=False) as exe: |
| 58 | + future_lst = [exe.submit(sleep_funct, 0.1) for _ in range(4)] |
| 59 | + self.assertEqual([f.done() for f in future_lst], [False, False, False, False]) |
| 60 | + self.assertEqual(len(exe), 4) |
| 61 | + self.assertEqual(exe.max_workers, 1) |
| 62 | + future_lst[0].result() |
| 63 | + exe.max_workers = 2 |
| 64 | + self.assertEqual(exe.max_workers, 2) |
| 65 | + self.assertEqual([f.done() for f in future_lst], [True, False, False, False]) |
| 66 | + self.assertEqual([f.result() for f in future_lst], [0.1, 0.1, 0.1, 0.1]) |
| 67 | + self.assertEqual([f.done() for f in future_lst], [True, True, True, True]) |
| 68 | + |
| 69 | + def test_no_block_allocation(self): |
| 70 | + with self.assertRaises(NotImplementedError): |
| 71 | + with SingleNodeExecutor(block_allocation=False, disable_dependencies=False) as exe: |
| 72 | + exe.max_workers = 2 |
| 73 | + with self.assertRaises(NotImplementedError): |
| 74 | + with SingleNodeExecutor(block_allocation=False, disable_dependencies=True) as exe: |
| 75 | + exe.max_workers = 2 |
| 76 | + |
| 77 | + def test_max_workers_stopped_executor(self): |
| 78 | + exe = SingleNodeExecutor(block_allocation=True) |
| 79 | + exe.shutdown(wait=True) |
| 80 | + self.assertIsNone(exe.max_workers) |
0 commit comments