Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 31 additions & 21 deletions src/executorlib/task_scheduler/interactive/blockallocation.py
Original file line number Diff line number Diff line change
Expand Up @@ -80,35 +80,37 @@ def __init__(
executor_kwargs["restart_limit"] = restart_limit
self._process_kwargs = executor_kwargs
self._max_workers = max_workers
self_id = random.getrandbits(128)
self._self_id = self_id
self._self_id = random.getrandbits(128)
_interrupt_bootup_dict[self._self_id] = False
alive_workers = [max_workers]
alive_workers_lock = Lock()
bootup_events = [Event() for _ in range(self._max_workers)]
bootup_events[0].set()
self._alive_workers = [max_workers]
self._alive_workers_lock = Lock()
self._bootup_events = [Event() for _ in range(self._max_workers)]
self._bootup_events[0].set()
self._set_process(
process=[
Thread(
target=_execute_multiple_tasks,
kwargs=executor_kwargs
| {
"worker_id": worker_id,
"stop_function": lambda: _interrupt_bootup_dict[self_id],
"bootup_event": bootup_events[worker_id],
"next_bootup_event": (
bootup_events[worker_id + 1]
if worker_id + 1 < self._max_workers
else None
),
"alive_workers": alive_workers,
"alive_workers_lock": alive_workers_lock,
},
kwargs=self._worker_kwargs(worker_id),
)
for worker_id in range(self._max_workers)
],
)

def _worker_kwargs(self, worker_id: int) -> dict:
self_id = self._self_id
return self._process_kwargs | {
"worker_id": worker_id,
"stop_function": lambda: _interrupt_bootup_dict[self_id],
"bootup_event": self._bootup_events[worker_id],
"next_bootup_event": (
self._bootup_events[worker_id + 1]
if worker_id + 1 < len(self._bootup_events)
else None
),
Comment thread
coderabbitai[bot] marked this conversation as resolved.
"alive_workers": self._alive_workers,
"alive_workers_lock": self._alive_workers_lock,
}

@property
def max_workers(self) -> int:
return self._max_workers
Expand All @@ -126,12 +128,20 @@ def max_workers(self, max_workers: int):
process for process in self._process if process.is_alive()
]
elif self._max_workers < max_workers:
old_max_workers = self._max_workers
self._bootup_events.extend(
Event() for _ in range(max_workers - old_max_workers)
)
for idx in range(old_max_workers, max_workers):
self._bootup_events[idx].set()
with self._alive_workers_lock:
self._alive_workers[0] += max_workers - old_max_workers
new_process_lst = [
Thread(
target=_execute_multiple_tasks,
kwargs=self._process_kwargs,
kwargs=self._worker_kwargs(worker_id),
)
for _ in range(max_workers - self._max_workers)
for worker_id in range(old_max_workers, max_workers)
]
for process_instance in new_process_lst:
process_instance.start()
Expand Down
49 changes: 46 additions & 3 deletions tests/unit/task_scheduler/interactive/test_blockallocation.py
Original file line number Diff line number Diff line change
@@ -1,11 +1,54 @@
import queue
import unittest
from threading import Lock
from concurrent.futures import Future
from threading import Event, Lock
from unittest.mock import patch

from executorlib.task_scheduler.interactive.blockallocation import _drain_dead_worker
from executorlib.task_scheduler.interactive.shared import task_done
from executorlib.standalone.interactive.communication import ExecutorlibSocketError
from executorlib.task_scheduler.interactive.blockallocation import (
BlockAllocationTaskScheduler,
_drain_dead_worker,
)


class TestBlockAllocationResize(unittest.TestCase):
def test_increase_workers_passes_worker_context(self):
scheduler = object.__new__(BlockAllocationTaskScheduler)
scheduler._future_queue = queue.Queue()
scheduler._process = []
scheduler._process_kwargs = {"future_queue": scheduler._future_queue}
scheduler._max_workers = 1
scheduler._self_id = 1
scheduler._alive_workers = [1]
scheduler._alive_workers_lock = Lock()
scheduler._bootup_events = [Event()]

class FakeThread:
instances = []

def __init__(self, target, kwargs):
self.target = target
self.kwargs = kwargs
self.started = False
self.instances.append(self)

def start(self):
self.started = True

with patch(
"executorlib.task_scheduler.interactive.blockallocation.Thread",
FakeThread,
):
scheduler.max_workers = 2

worker = FakeThread.instances[-1]
self.assertEqual(worker.kwargs["worker_id"], 1)
self.assertIn("stop_function", worker.kwargs)
self.assertIn("bootup_event", worker.kwargs)
self.assertIn("next_bootup_event", worker.kwargs)
self.assertIs(worker.kwargs["alive_workers"], scheduler._alive_workers)
self.assertTrue(worker.started)
self.assertEqual(scheduler._alive_workers[0], 2)


class TestDrainDeadWorker(unittest.TestCase):
Expand Down
Loading