Skip to content
Open
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
2 changes: 1 addition & 1 deletion src/executorlib/standalone/interactive/communication.py
Original file line number Diff line number Diff line change
Expand Up @@ -152,7 +152,7 @@ def shutdown(self, wait: bool = True):
if self._spawner.poll():
result = self.send_and_receive_dict(
input_dict={"shutdown": True, "wait": wait}
)["result"]
).get("result")
self._spawner.shutdown(wait=wait)
self._reset_socket()
return result
Expand Down
51 changes: 30 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 on lines +99 to +109

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

Release the first added worker from bootup.

_worker_kwargs() snapshots the original tail worker's next_bootup_event as None. After resize, the first added worker waits on the new unset event, but no existing worker can signal it. The added worker and later added workers remain blocked in bootup_event.wait().

Signal the first new bootup event when resizing. If strict worker-ID boot order must also apply during concurrent startup, use a resizable handoff instead of a static next_bootup_event.

Proposed fix
                 self._bootup_events.extend(
                     Event() for _ in range(max_workers - old_max_workers)
                 )
+                self._bootup_events[old_max_workers].set()
                 self._alive_workers[0] += max_workers - old_max_workers

Also applies to: 130-140

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@src/executorlib/task_scheduler/interactive/blockallocation.py` around lines
99 - 108, Update the resize path and _worker_kwargs so the first newly added
worker’s bootup event is explicitly released, while preserving sequential
handoff for later workers. Ensure concurrent startup cannot leave added workers
waiting on a statically captured next_bootup_event; use the existing
resize/startup synchronization mechanism or a resizable handoff, and apply the
same correction to the corresponding logic around the later referenced block.

"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,19 @@ 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()
self._alive_workers[0] += max_workers - old_max_workers

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🩺 Stability & Availability | 🟠 Major | ⚡ Quick win

Synchronize the live-worker count update.

_drain_dead_worker() decrements self._alive_workers[0] under self._alive_workers_lock. A worker can fail while this resize increments the same value. An unsynchronized increment can lose either update and make the scheduler treat live workers as dead.

Proposed fix
-                self._alive_workers[0] += max_workers - old_max_workers
+                with self._alive_workers_lock:
+                    self._alive_workers[0] += max_workers - old_max_workers
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
self._alive_workers[0] += max_workers - old_max_workers
with self._alive_workers_lock:
self._alive_workers[0] += max_workers - old_max_workers
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@src/executorlib/task_scheduler/interactive/blockallocation.py` at line 134,
Protect the live-worker count increment in the resize logic with
self._alive_workers_lock, matching the locking used by _drain_dead_worker().
Keep the adjustment of self._alive_workers[0] atomic with respect to concurrent
worker-failure decrements.

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