diff --git a/CHANGELOG.rst b/CHANGELOG.rst index a7af6272..4c2c37ad 100644 --- a/CHANGELOG.rst +++ b/CHANGELOG.rst @@ -3,6 +3,8 @@ * `#380 `__: Add support for Python 3.13 and 3.14, and drop EOL 3.8 and 3.9. +* `#429 `__: ``safe_terminate`` now waits for a terminated worker once more after the kill attempt, so a released worker is joined instead of abandoned, and returns whether all workers finished within the bounds instead of discarding the wait result. ``Group.terminate()`` now propagates that result to its callers. + 2.1.2 (2025-11-11) ------------------ diff --git a/news/429.bugfix.rst b/news/429.bugfix.rst new file mode 100644 index 00000000..78d1e974 --- /dev/null +++ b/news/429.bugfix.rst @@ -0,0 +1 @@ +``safe_terminate`` now waits for a terminated worker once more after the kill attempt, so a released worker is joined instead of abandoned, and returns whether all workers finished within the bounds. ``Group.terminate()`` propagates that result to its callers. diff --git a/src/execnet/multi.py b/src/execnet/multi.py index 4dbf8b89..4d6df88e 100644 --- a/src/execnet/multi.py +++ b/src/execnet/multi.py @@ -208,7 +208,7 @@ def _cleanup_atexit(self) -> None: trace(f"=== atexit cleanup {self!r} ===") self.terminate(timeout=1.0) - def terminate(self, timeout: float | None = None) -> None: + def terminate(self, timeout: float | None = None) -> bool: """Trigger exit of member gateways and wait for termination of member gateways and associated subprocesses. @@ -217,7 +217,12 @@ def terminate(self, timeout: float | None = None) -> None: Timeout defaults to None meaning open-ended waiting and no kill attempts. + + Returns whether all member gateways finished within the bounds; + ``False`` means some gateway is still running and its thread was + abandoned. """ + all_finished = True while self: vias: set[str] = set() for gw in self: @@ -235,15 +240,19 @@ def kill(gw: Gateway) -> None: trace("Gateways did not come down after timeout: %r" % gw) gw._io.kill() - safe_terminate( - self.execmodel, - timeout, - [ - (partial(join_wait, gw), partial(kill, gw)) - for gw in self._gateways_to_join - ], + all_finished = ( + safe_terminate( + self.execmodel, + timeout, + [ + (partial(join_wait, gw), partial(kill, gw)) + for gw in self._gateways_to_join + ], + ) + and all_finished ) self._gateways_to_join[:] = [] + return all_finished def remote_exec( self, @@ -337,12 +346,17 @@ def safe_terminate( execmodel: ExecModel, timeout: float | None, list_of_paired_functions: Sequence[TermKillPair], -) -> None: +) -> bool: """Run terminate/kill pairs in parallel with a hard wait bound. - Each termfunc is given ``timeout``. If it does not finish, killfunc runs. + Each termfunc is given ``timeout``. If it does not finish, killfunc runs, + after which the termfunc gets one more ``timeout`` to finish, so a kill + that releases it lets the worker be joined rather than abandoned. Waiting for the worker pool is also bounded so a stuck kill cannot hang the caller forever (see issues #43 / #221). + + Returns whether all workers finished within the bounds; ``False`` means + some termfunc is still running and its thread was abandoned. """ workerpool = WorkerPool(execmodel) @@ -350,15 +364,24 @@ def termkill(termfunc: TermKillFunc, killfunc: TermKillFunc) -> None: termreply = workerpool.spawn(termfunc) try: termreply.get(timeout=timeout) + return except OSError: killfunc() + # The kill should release the termfunc; observe it finishing so the + # worker is joined instead of abandoned (issue #429). A kill that + # does not release it is reported through the waitall result below. + try: + termreply.waitfinish(timeout=timeout) + except OSError: + pass replylist = [ workerpool.spawn(termkill, termfunc, killfunc) for termfunc, killfunc in list_of_paired_functions ] - # Allow term timeout plus a kill attempt; never block indefinitely. - wait_timeout = None if timeout is None else timeout * 2 + # Allow term timeout, a kill attempt and one more term timeout after the + # kill; never block indefinitely. + wait_timeout = None if timeout is None else timeout * 3 for reply in replylist: try: reply.waitfinish(timeout=wait_timeout) @@ -366,7 +389,7 @@ def termkill(termfunc: TermKillFunc, killfunc: TermKillFunc) -> None: # termkill still running (typically stuck in killfunc). continue reply.get() # propagate worker exceptions, if any - workerpool.waitall(timeout=wait_timeout) + return workerpool.waitall(timeout=wait_timeout) default_group = Group() diff --git a/testing/test_multi.py b/testing/test_multi.py index 12e0ed3d..528d348c 100644 --- a/testing/test_multi.py +++ b/testing/test_multi.py @@ -316,3 +316,61 @@ def kill_ok() -> None: assert kill_started.is_set() assert other_killed == [1] release_kill.set() + + +@pytest.mark.timeout(10) +def test_safe_terminate_reports_kill_that_ignores( + execmodel: ExecModel, +) -> None: + """Regression for #429: a kill that does not release the termfunc is + reported through the return value instead of being silently discarded.""" + if execmodel.backend not in ("thread", "main_thread_only"): + pytest.xfail( + "execution model %r does not support task count" % execmodel.backend + ) + entered = execmodel.Event() + release = execmodel.Event() + + def term() -> None: + entered.set() + release.wait() + + def kill() -> None: + pass # ignores the kill, leaving the termfunc running + + try: + result = safe_terminate(execmodel, 0.2, [(term, kill)]) + finally: + release.set() + + assert entered.is_set() + assert result is False + + +@pytest.mark.timeout(10) +def test_safe_terminate_joins_when_kill_releases( + execmodel: ExecModel, +) -> None: + """Regression for #429: a working kill lets every worker be observed to + finish, so safe_terminate returns success with no abandoned threads.""" + if execmodel.backend not in ("thread", "main_thread_only"): + pytest.xfail( + "execution model %r does not support task count" % execmodel.backend + ) + entered = execmodel.Event() + release = execmodel.Event() + finished = execmodel.Event() + + def term() -> None: + entered.set() + release.wait() + finished.set() + + def kill() -> None: + release.set() + + result = safe_terminate(execmodel, 1, [(term, kill)]) + + assert entered.is_set() + assert result is True + assert finished.is_set() # joined before returning, no stragglers