Skip to content

ProcessPoolExecutor: after a worker dies, shutdown(wait=False) and submit() block until every other worker has exited #158413

Description

@kvnlng

Bug report

Bug description:

When a ProcessPoolExecutor worker dies abruptly, the executor manager thread runs terminate_broken(). That holds _shutdown_lock while it sends every remaining worker SIGTERM, sends SIGTERM again in _join_executor_internals(broken=True), and then calls p.join() on each worker with no timeout. submit() takes _shutdown_lock first, as does shutdown() whatever wait is, and so does 3.14's kill_workers() (the repro calls it; terminate_workers() goes through the same _force_shutdown()). So if a remaining worker does not exit on SIGTERM, all of these calls block until it does, which can be never: once idle, the worker waits in call_queue.get() and never sees EOF there, because it holds a copy of the pipe's write end itself.

Workers outlive SIGTERM in ordinary programs. A SIGTERM handler that the main script installs at module level runs again in every spawn worker. coverage.py's sigterm = true handler can deadlock, or raise into the running task instead of ending the process (coveragepy/coveragepy#2310). Even a handler that calls sys.exit() is not enough when the second SIGTERM arrives while the task is still unwinding, for example in a finally block: _process_worker catches that SystemExit as the task's result and goes back to call_queue.get().

"""After a ProcessPoolExecutor worker dies, shutdown(wait=False) and submit()
block for as long as another worker outlives SIGTERM.

    python repro.py ignore    # workers' SIGTERM handler does nothing
    python repro.py exit      # workers' SIGTERM handler calls sys.exit(0)
    python repro.py default   # control: default SIGTERM action

Bounded: after BOUND seconds the script SIGKILLs the workers itself.
"""
import concurrent.futures
import multiprocessing
import os
import signal
import sys
import threading
import time

BOUND = 10.0


def install_handler(mode):
    if mode == "ignore":
        signal.signal(signal.SIGTERM, lambda signum, frame: None)
    elif mode == "exit":
        signal.signal(signal.SIGTERM, lambda signum, frame: sys.exit(0))


def work(seconds):
    try:
        time.sleep(seconds)
    finally:
        time.sleep(0.05)  # a task's cleanup, e.g. closing a file
    return os.getpid()


def die(delay):
    time.sleep(delay)
    os._exit(1)


def in_thread(fn):
    outcome = []

    def run():
        try:
            fn()
            outcome.append("returned")
        except Exception as exc:
            outcome.append(f"raised {type(exc).__name__}")

    thread = threading.Thread(target=run, daemon=True)
    thread.start()
    return thread, outcome


def main():
    mode = sys.argv[1]
    ex = concurrent.futures.ProcessPoolExecutor(
        max_workers=2, mp_context=multiprocessing.get_context("spawn"),
        initializer=install_handler, initargs=(mode,))
    assert len({f.result() for f in [ex.submit(work, 0.5) for _ in range(2)]}) == 2
    workers = list(ex._processes.values())
    manager = ex._executor_manager_thread

    # One worker is running a task when the other one dies.
    busy = ex.submit(work, 1.0)
    ex.submit(die, 0.2)
    try:
        busy.result()
    except concurrent.futures.process.BrokenProcessPool:
        print("pool is broken")

    calls = {"submit()": in_thread(lambda: ex.submit(work, 0)),
             "shutdown(wait=False)": in_thread(lambda: ex.shutdown(wait=False))}
    if hasattr(ex, "kill_workers"):  # 3.14+
        calls["kill_workers()"] = in_thread(ex.kill_workers)
    deadline = time.monotonic() + BOUND
    for thread, _ in calls.values():
        thread.join(max(0.0, deadline - time.monotonic()))

    blocked = [name for name, (thread, _) in calls.items() if thread.is_alive()]
    for name, (_, outcome) in calls.items():
        print(f"{name}: {outcome[0] if outcome else f'still blocked after {BOUND:.0f} s'}")
    if blocked:
        frame = sys._current_frames().get(manager.ident)
        stack = []
        while frame is not None:
            stack.append(frame.f_code.co_name)
            frame = frame.f_back
        print("executor manager thread:", " <- ".join(stack))
        alive = [p for p in workers if p.is_alive()]
        print("workers still alive:", [p.pid for p in alive])
        for p in alive:
            p.kill()
        for thread, _ in calls.values():
            thread.join(BOUND)
        for name in blocked:
            outcome = calls[name][1]
            print(f"after SIGKILL: {name} {outcome[0] if outcome else 'STILL BLOCKED'}")
    for child in multiprocessing.active_children():
        child.kill()
    if any(thread.is_alive() for thread, _ in calls.values()):
        os._exit(2)  # the atexit hook would join the manager thread
    sys.exit(1 if blocked else 0)


if __name__ == "__main__":
    main()

python repro.py ignore on 3.14.7:

pool is broken
submit(): still blocked after 10 s
shutdown(wait=False): still blocked after 10 s
kill_workers(): still blocked after 10 s
executor manager thread: poll <- wait <- join <- _join_executor_internals <- _terminate_broken <- terminate_broken <- run <- _bootstrap_inner <- _bootstrap
workers still alive: [43687]
after SIGKILL: submit() raised BrokenProcessPool
after SIGKILL: shutdown(wait=False) returned
after SIGKILL: kill_workers() returned

python repro.py exit prints the same. 3.12 prints the same without the kill_workers() lines. python repro.py default finishes in about a second, with submit(): raised BrokenProcessPool, shutdown(wait=False): returned and kill_workers(): returned.

Expected: shutdown(wait=False) returns at once, submit() raises BrokenProcessPool at once, and kill_workers() kills the worker that is left. Actual: all three block until something outside the executor kills that worker.

ignore and exit blocked in every run and default in none: 20 runs of each on CPython 3.12.14 and on the 3.14.7 free-threaded build (PYTHON_GIL=0), and 10 of each on the standard 3.14.7 build, all on macOS 27.0 (arm64).

The code, in Lib/concurrent/futures/process.py (line numbers for v3.12.14 / v3.14.7 / main at c8ea867):

A possible direction: bound the joins in _join_executor_internals(broken=True) (escalating to kill()), or don't hold _shutdown_lock across them, so that shutdown(wait=False), submit() and kill_workers() never wait on a worker.

CPython versions tested on:

3.12, 3.14

Operating systems tested on:

macOS

Linked PRs

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    stdlibStandard Library Python modules in the Lib/ directorytype-bugAn unexpected behavior, bug, or error

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions