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: 2 additions & 0 deletions CHANGELOG.rst
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@

* `#380 <https://github.com/pytest-dev/execnet/pull/380>`__: Add support for Python 3.13 and 3.14, and drop EOL 3.8 and 3.9.

* `#429 <https://github.com/pytest-dev/execnet/issues/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)
------------------

Expand Down
19 changes: 16 additions & 3 deletions src/execnet/multi.py
Original file line number Diff line number Diff line change
Expand Up @@ -337,21 +337,34 @@ 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)

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)
Expand All @@ -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()
Expand Down
58 changes: 58 additions & 0 deletions testing/test_multi.py
Original file line number Diff line number Diff line change
Expand Up @@ -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