Skip to content

Commit 11c7bcb

Browse files
authored
Add test for exception handler in _update_waiting_task batched path
1 parent 0a12ed7 commit 11c7bcb

2 files changed

Lines changed: 60 additions & 16 deletions

File tree

src/executorlib/_version.py

Lines changed: 18 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1,21 +1,24 @@
1-
# file generated by setuptools-scm
1+
# file generated by vcs-versioning
22
# don't change, don't track in version control
3+
from __future__ import annotations
34

4-
__all__ = ["__version__", "__version_tuple__", "version", "version_tuple"]
5-
6-
TYPE_CHECKING = False
7-
if TYPE_CHECKING:
8-
from typing import Tuple
9-
from typing import Union
10-
11-
VERSION_TUPLE = Tuple[Union[int, str], ...]
12-
else:
13-
VERSION_TUPLE = object
5+
__all__ = [
6+
"__version__",
7+
"__version_tuple__",
8+
"version",
9+
"version_tuple",
10+
"__commit_id__",
11+
"commit_id",
12+
]
1413

1514
version: str
1615
__version__: str
17-
__version_tuple__: VERSION_TUPLE
18-
version_tuple: VERSION_TUPLE
16+
__version_tuple__: tuple[int | str, ...]
17+
version_tuple: tuple[int | str, ...]
18+
commit_id: str | None
19+
__commit_id__: str | None
20+
21+
__version__ = version = '0.1.dev2+g0a12ed789.d20260610'
22+
__version_tuple__ = version_tuple = (0, 1, 'dev2', 'g0a12ed789.d20260610')
1923

20-
__version__ = version = "0.0.1"
21-
__version_tuple__ = version_tuple = (0, 0, 1)
24+
__commit_id__ = commit_id = None

tests/unit/executor/test_single_dependencies.py

Lines changed: 42 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,10 +3,14 @@
33
from time import sleep, time
44
from queue import Queue
55
from threading import Thread
6+
from unittest.mock import MagicMock
67

78
from executorlib import SingleNodeExecutor
89
from executorlib.executor.single import create_single_node_executor
9-
from executorlib.task_scheduler.interactive.dependency import _execute_tasks_with_dependencies
10+
from executorlib.task_scheduler.interactive.dependency import (
11+
_execute_tasks_with_dependencies,
12+
_update_waiting_task,
13+
)
1014
from executorlib.standalone.serialize import cloudpickle_register
1115
from executorlib.standalone.interactive.spawner import MpiExecSpawner
1216

@@ -300,6 +304,43 @@ def test_future_input_dict(self):
300304
)
301305
self.assertEqual(fs.result()["a"], 4)
302306

307+
def test_update_waiting_task_batched_exception(self):
308+
"""_update_waiting_task catches exceptions from batched_futures and sets them on the batch future."""
309+
executor_queue = Queue()
310+
batch_future = Future()
311+
312+
# A mock skip_lst future: done(), exception() returns None (passes get_exception_lst),
313+
# but result() raises -- triggering the except block in _update_waiting_task.
314+
mock_skip_future = MagicMock()
315+
mock_skip_future.done.return_value = True
316+
mock_skip_future.exception.return_value = None
317+
mock_skip_future.result.side_effect = RuntimeError("unexpected skip error")
318+
319+
task_dict = {
320+
"fn": "batched",
321+
"args": (),
322+
"kwargs": {
323+
"lst": [],
324+
"n": 3,
325+
"skip_lst": [mock_skip_future],
326+
},
327+
"future": batch_future,
328+
"future_lst": [mock_skip_future],
329+
"resource_dict": {},
330+
}
331+
332+
result_lst = _update_waiting_task(
333+
wait_lst=[task_dict],
334+
executor_queue=executor_queue,
335+
refresh_rate=0.0,
336+
)
337+
338+
# The batch future must have the exception propagated (not crashed the scheduler)
339+
self.assertTrue(batch_future.done())
340+
self.assertIsInstance(batch_future.exception(), RuntimeError)
341+
# The failed task is consumed (not re-queued in the wait list)
342+
self.assertEqual(len(result_lst), 0)
343+
303344

304345
class TestExecutorErrors(unittest.TestCase):
305346
def test_block_allocation_false_one_worker(self):

0 commit comments

Comments
 (0)