From c9cb849aeacdbb032175169ea609a959d53eee1f Mon Sep 17 00:00:00 2001 From: inchang-ing <197932532+inchang-ing@users.noreply.github.com> Date: Mon, 28 Sep 2026 20:07:56 +0800 Subject: [PATCH 1/3] fix: observe terminated workers and report safe_terminate outcome After the kill attempt, wait once more (bounded by the same timeout) for the termfunc worker, so a kill that releases it lets the worker be joined instead of abandoned while still running. Threads spawned via _thread.start_new_thread are never joined by threading._shutdown, so abandoned workers could outlive safe_terminate arbitrarily. safe_terminate now also returns the WorkerPool.waitall result instead of discarding it, letting Group.terminate and other callers tell whether all workers actually finished within the bounds. Fixes #429 --- CHANGELOG.rst | 2 ++ src/execnet/multi.py | 19 +++++++++++--- testing/test_multi.py | 58 +++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 76 insertions(+), 3 deletions(-) diff --git a/CHANGELOG.rst b/CHANGELOG.rst index a7af6272..8f13df8e 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. + 2.1.2 (2025-11-11) ------------------ diff --git a/src/execnet/multi.py b/src/execnet/multi.py index 4dbf8b89..7f220bb4 100644 --- a/src/execnet/multi.py +++ b/src/execnet/multi.py @@ -337,12 +337,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,8 +355,16 @@ 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) @@ -366,7 +379,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 From fc5b814e0215bc9ed76e353fc3bc597dd7b95c18 Mon Sep 17 00:00:00 2001 From: inchang-ing <197932532+inchang-ing@users.noreply.github.com> Date: Thu, 8 Oct 2026 01:53:07 +0800 Subject: [PATCH 2/3] fix: propagate the safe_terminate outcome and widen the outer wait Group.terminate() now returns whether all member gateways finished, so pytest-xdist-style callers can use the result. The outer wait budget grows to timeout * 3 because termkill can now take timeout + kill + timeout. --- CHANGELOG.rst | 2 +- src/execnet/multi.py | 30 ++++++++++++++++++++---------- 2 files changed, 21 insertions(+), 11 deletions(-) diff --git a/CHANGELOG.rst b/CHANGELOG.rst index 8f13df8e..4c2c37ad 100644 --- a/CHANGELOG.rst +++ b/CHANGELOG.rst @@ -3,7 +3,7 @@ * `#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. +* `#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/src/execnet/multi.py b/src/execnet/multi.py index 7f220bb4..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, @@ -370,8 +379,9 @@ def termkill(termfunc: TermKillFunc, killfunc: TermKillFunc) -> None: 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) From c24f168f4569a880cbde302b205083e05f370e63 Mon Sep 17 00:00:00 2001 From: inchang-ing <197932532+inchang-ing@users.noreply.github.com> Date: Thu, 8 Oct 2026 07:52:39 +0800 Subject: [PATCH 3/3] chore: add the news fragment required by the changelog check The Changelog entry check looks for news/..rst, which this branch was missing; CHANGELOG.rst is generated from these fragments. --- news/429.bugfix.rst | 1 + 1 file changed, 1 insertion(+) create mode 100644 news/429.bugfix.rst 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.